diff --git a/agent/inline_tool_executors.py b/agent/inline_tool_executors.py new file mode 100644 index 0000000000..da5c8b6c81 --- /dev/null +++ b/agent/inline_tool_executors.py @@ -0,0 +1,292 @@ +"""Agent-level ("inline") tool executors shared by the sequential and concurrent tool paths. + +These tools need live ``AIAgent`` state (stores, callbacks, session DB) and therefore +bypass the tool registry. Each executor is ``fn(agent, args, ctx) -> result``; the +table replaces two hand-maintained if/elif chains (``invoke_tool`` and +``execute_tool_calls_sequential``) that had drifted apart. Tool modules are imported +lazily inside the bodies so ``patch("tools.x.y")`` in tests keeps working. +""" + +from __future__ import annotations + +import json +from dataclasses import dataclass +from typing import Any, Callable, Dict, Optional + + +def tool_hook_ids(agent, effective_task_id: str, tool_call_id: Optional[str]) -> Dict[str, str]: + """Identity kwargs every tool hook/middleware call carries (all coerced to ``""``).""" + return { + "task_id": effective_task_id or "", + "session_id": getattr(agent, "session_id", "") or "", + "tool_call_id": tool_call_id or "", + "turn_id": getattr(agent, "_current_turn_id", "") or "", + "api_request_id": getattr(agent, "_current_api_request_id", "") or "", + } + + +def emit_terminal_post_tool_call( + agent, + *, + function_name: str, + function_args: dict, + result: Any, + effective_task_id: str, + tool_call_id: Optional[str], + duration_ms: int = 0, + status: Optional[str] = None, + error_type: Optional[str] = None, + error_message: Optional[str] = None, + middleware_trace: Optional[list] = None, +) -> None: + """Emit the one terminal ``post_tool_call`` hook for a tool_call_id (best-effort).""" + try: + from model_tools import _emit_post_tool_call_hook + _emit_post_tool_call_hook( + function_name=function_name, + function_args=function_args, + result=result, + **tool_hook_ids(agent, effective_task_id, tool_call_id), + duration_ms=duration_ms, + status=status, + error_type=error_type, + error_message=error_message, + middleware_trace=list(middleware_trace or []), + ) + except Exception: + pass + + +@dataclass +class InlineToolContext: + """Per-call state an inline executor may need beyond its arguments.""" + + effective_task_id: str + tool_call_id: Optional[str] = None + messages: Optional[list] = None + + +def _todo_list(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.todo_tool import todo_tool as _todo_tool + + return _todo_tool( + todos=args.get("todos"), + merge=args.get("merge", False), + store=agent._todo_store, + ) + + +def _message_agent(agent, args: dict, ctx: InlineToolContext) -> Any: + # Bot Mode teammate DM is injected, not registered: only a canonical Bot + # Chat session carries the schema, and the tool re-gates on the title. + from tools.bot_mode_dm import message_agent_tool as _message_agent_tool + + return _message_agent_tool( + target=args.get("target", ""), + message=args.get("message", ""), + task_id=ctx.effective_task_id, + agent=agent, + ) + + +def _session_search(agent, args: dict, ctx: InlineToolContext) -> Any: + session_db = agent._get_session_db_for_recall() + if not session_db: + from hermes_state import format_session_db_unavailable + + return json.dumps({"success": False, "error": format_session_db_unavailable()}) + from tools.session_search_tool import session_search as _session_search_tool + + return _session_search_tool( + query=args.get("query", ""), + role_filter=args.get("role_filter"), + limit=args.get("limit", 3), + session_id=args.get("session_id"), + around_message_id=args.get("around_message_id"), + window=args.get("window", 5), + sort=args.get("sort"), + detail=args.get("detail", "adaptive"), + db=session_db, + current_session_id=agent.session_id, + ) + + +def _memory(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.memory_tool import memory_tool as _memory_tool + + result = _memory_tool( + action=args.get("action"), + target=args.get("target", "memory"), + content=args.get("content"), + old_text=args.get("old_text"), + operations=args.get("operations"), + store=agent._memory_store, + ) + # Mirror built-in memory writes to external providers; gating lives in + # MemoryManager.notify_memory_tool_write. + if agent._memory_manager: + agent._memory_manager.notify_memory_tool_write( + result, + args, + build_metadata=lambda: agent._build_memory_write_metadata( + task_id=ctx.effective_task_id, + tool_call_id=ctx.tool_call_id, + ), + ) + return result + + +def _clarify(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.clarify_tool import clarify_tool as _clarify_tool + + return _clarify_tool( + question=args.get("question", ""), + choices=args.get("choices"), + multi_select=args.get("multi_select", False), + questions=args.get("questions"), + callback=agent.clarify_callback, + ) + + +def _read_terminal(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.read_terminal_tool import read_terminal_tool as _read_terminal_tool + + return _read_terminal_tool( + start_line=args.get("start_line"), + count=args.get("count"), + callback=getattr(agent, "read_terminal_callback", None), + ) + + +def _desktop_preview(agent, args: dict, ctx: InlineToolContext) -> Any: + # action=read needs the GUI callback (agent-level); open/close go through the + # registry handler like any other tool. + if (args.get("action") or "").strip() == "read": + from tools.read_preview_tool import read_preview_tool as _read_preview_tool + + return _read_preview_tool( + start=args.get("start"), + count=args.get("count"), + callback=getattr(agent, "read_preview_callback", None), + ) + from tools.preview_tool import _handle_preview + + return _handle_preview(args) + + +def _drive_preview(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.drive_preview_tool import drive_preview_tool as _drive_preview_tool + + return _drive_preview_tool( + action=args.get("action", ""), + ref=args.get("ref"), + selector=args.get("selector"), + text=args.get("text"), + key=args.get("key"), + submit=args.get("submit"), + amount=args.get("amount"), + to=args.get("to"), + limit=args.get("max"), + callback=getattr(agent, "drive_preview_callback", None), + ) + + +def _annotate_preview(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.annotate_preview_tool import annotate_preview_tool as _annotate_preview_tool + + return _annotate_preview_tool( + action=args.get("action", "add"), + ref=args.get("ref"), + selector=args.get("selector"), + label=args.get("label"), + callback=getattr(agent, "drive_preview_callback", None), + ) + + +def _read_window_below(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.read_window_tool import read_window_below_tool as _read_window_below_tool + + return _read_window_below_tool( + callback=getattr(agent, "read_window_below_callback", None), + ) + + +def _gui_tour(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.tour_tool import tour_tool as _tour_tool + + return _tour_tool( + action=args.get("action", ""), + surface=args.get("surface"), + selector=args.get("selector"), + title=args.get("title"), + text=args.get("text"), + side=args.get("side"), + steps=args.get("steps"), + step_index=args.get("step_index"), + callback=getattr(agent, "tour_callback", None), + ) + + +def _setup_mcp(agent, args: dict, ctx: InlineToolContext) -> Any: + from tools.setup_mcp_tool import setup_mcp_tool as _setup_mcp_tool + + return _setup_mcp_tool( + server=args.get("server", ""), + action=args.get("action", "install"), + reason=args.get("reason", ""), + callback=getattr(agent, "setup_mcp_callback", None), + ) + + +def _delegate_task(agent, args: dict, ctx: InlineToolContext) -> Any: + return agent._dispatch_delegate_task(args) + + +InlineToolExecutor = Callable[[Any, dict, InlineToolContext], Any] + +# Order is the historical if/elif order of ``execute_tool_calls_sequential``. +INLINE_TOOL_EXECUTORS: Dict[str, InlineToolExecutor] = { + "todo_list": _todo_list, + "message_agent": _message_agent, + "session_search": _session_search, + "memory": _memory, + "clarify": _clarify, + "read_terminal": _read_terminal, + "desktop_preview": _desktop_preview, + "drive_preview": _drive_preview, + "annotate_preview": _annotate_preview, + "read_window_below": _read_window_below, + "gui_tour": _gui_tour, + "setup_mcp": _setup_mcp, + "delegate_task": _delegate_task, +} + +# ``invoke_tool`` (concurrent path) historically consulted the memory manager right +# after these three names and before the remaining inline tools; it never handled +# ``message_agent`` inline (that name falls through to the registry there). +INVOKE_TOOL_PRE_MEMORY_MANAGER_NAMES = frozenset({"todo_list", "session_search", "memory"}) + + +def memory_manager_executor(function_name: str) -> InlineToolExecutor: + """Executor routing ``function_name`` through ``agent._memory_manager``.""" + + def _run(agent, args: dict, ctx: InlineToolContext) -> Any: + return agent._memory_manager.handle_tool_call(function_name, args) + + return _run + + +def resolve_invoke_tool_executor(agent, function_name: str) -> Optional[InlineToolExecutor]: + """Inline executor for ``invoke_tool`` (concurrent path), or None for registry dispatch. + + Preserves the historical precedence: todo_list/session_search/memory, then memory + manager tools, then the remaining inline tools (``message_agent`` excluded). + """ + if function_name in INVOKE_TOOL_PRE_MEMORY_MANAGER_NAMES: + return INLINE_TOOL_EXECUTORS[function_name] + memory_manager = agent._memory_manager + if memory_manager and memory_manager.has_tool(function_name): + return memory_manager_executor(function_name) + if function_name == "message_agent": + return None + return INLINE_TOOL_EXECUTORS.get(function_name) diff --git a/agent/tool_executor.py b/agent/tool_executor.py index f1a04c3718..af1b3c676b 100644 --- a/agent/tool_executor.py +++ b/agent/tool_executor.py @@ -1,13 +1,7 @@ -"""Tool-call execution — sequential and concurrent dispatch. +"""Tool-call execution: sequential and concurrent dispatch, extracted from AIAgent. -Both AIAgent methods (``_execute_tool_calls_sequential`` and -``_execute_tool_calls_concurrent``) live here as module-level -functions that take the parent ``AIAgent`` as their first argument. - -``run_agent`` keeps thin wrappers so existing call sites work; tests -that patch ``run_agent._set_interrupt`` are honored because the -extracted functions reach back through the ``run_agent`` module via -``_ra()`` for that symbol. +Functions take the parent ``AIAgent`` first; ``run_agent`` keeps thin wrappers, and +tests that patch ``run_agent._set_interrupt`` still work because we reach it via ``_ra()``. """ from __future__ import annotations @@ -33,6 +27,12 @@ from agent.display import ( _detect_tool_failure, ) from agent.message_sanitization import coalesce_tool_call_id +from agent.inline_tool_executors import ( + INLINE_TOOL_EXECUTORS, + InlineToolContext, + emit_terminal_post_tool_call, + tool_hook_ids, +) from agent.tool_dispatch_helpers import ( _NEVER_PARALLEL_TOOLS, _is_destructive_command, @@ -62,12 +62,9 @@ def _pairing_tool_call_id(tool_call: Any) -> str: def _record_persisted_path_for_stub(agent, tool_call_id: str, function_result) -> None: - """Tell the stall guards where a persisted result's full content lives. + """Record the spillover file path so a later result-reference stub can't dangle. - When a large result is spilled to disk ( preview), a - later result-reference stub pointing at that first occurrence must carry - the spillover file path so the reference can't dangle. Best-effort: never - lets bookkeeping break tool execution. + Best-effort: bookkeeping never breaks tool execution. """ try: if not isinstance(function_result, str): @@ -90,10 +87,8 @@ def _ensure_file_checkpoint( if not file_path: return - # File tools resolve relative paths against the task's live/session cwd, - # which can differ from the Hermes process cwd (notably in Docker). Resolve - # through that same path pipeline before asking the checkpoint manager to - # discover the project root. + # File tools resolve relative paths against the task's live cwd (differs from the + # process cwd in Docker); resolve the same way before locating the project root. from tools.file_tools import _resolve_path_for_task resolved_path = _resolve_path_for_task(file_path, effective_task_id or "default") @@ -104,17 +99,13 @@ def _ensure_file_checkpoint( def _budget_for_agent(agent) -> BudgetConfig: """Resolve a tool-result BudgetConfig scaled to the agent's context window. - Large-context models keep the historical 100K/200K char defaults; small - models (e.g. a 65K-token local model switched into mid-session) get a budget - proportional to their window so a single large tool result can't push the - request past the model's limit (#23767). Falls back to the default budget - when the context length isn't resolvable. + Small-context models get a proportional budget so one large result can't overflow + the request (#23767); falls back to the default when context length is unknown. """ try: ctx = getattr(getattr(agent, "context_compressor", None), "context_length", None) - # budget_for_context_window(None) (rather than DEFAULT_BUDGET) so the - # config-driven MCP threshold override still applies when the context - # length isn't resolvable. + # budget_for_context_window(None), not DEFAULT_BUDGET, so the MCP threshold + # override still applies when the context length isn't resolvable. return budget_for_context_window(int(ctx) if ctx else None) except Exception: return DEFAULT_BUDGET @@ -126,40 +117,26 @@ _DEFAULT_IMAGE_PARALLEL_REQUESTS = 4 # Generous ceiling for slow-but-valid tool work (large page fetches, slow # remote backends) so the batch guard does not preempt a legitimate attempt. _DEFAULT_CONCURRENT_TOOL_TIMEOUT_S = 420.0 -# Upper bound a concurrent worker will wait at the start-order gate for all -# earlier-ordered tools to advance before proceeding out of order. Long enough -# to cover slow-but-legitimate authorization (e.g. an approval round-trip), -# short enough that one wedged dispatch cannot starve the batch forever. +# Start-order gate wait bound: long enough for an approval round-trip, short enough +# that one wedged dispatch cannot starve the batch. _START_ORDER_GATE_TIMEOUT_S = 120.0 -# Fallback bound a concurrent worker will wait for the authorization gate's -# serialization lock before running its prompt unserialized. The effective -# bound is derived from ``approvals.timeout`` plus a margin (see -# _authorization_gate_lock_timeout): a legitimate holder is at worst a human -# answering an approval prompt, which self-terminates at approvals.timeout — -# so a holder that overstays it is wedged and must not starve the batch. +# Fallback authorization-gate lock bound; the effective bound derives from +# approvals.timeout (see _authorization_gate_lock_timeout) since overstaying it means wedged. _AUTHORIZATION_GATE_LOCK_TIMEOUT_S = 360.0 def _authorization_gate_lock_timeout() -> float: """Bound for the authorization serialization lock: approval timeout + margin. - Delegates to ``tools.approval.human_wait_ceiling`` — the same bound that - clamps a human-wait window's deadline contribution — so the two can't - drift. Long enough that serialization is never broken while a legitimate - approval prompt is still answerable; short enough that a wedged holder - (hanging ``pre_tool_call`` plugin, dead approval client) cannot park other - workers forever (#79719). Resolved once per gate (per batch), so a - mid-process ``approvals.timeout`` change applies from the next batch. + Delegates to ``tools.approval.human_wait_ceiling`` so the two bounds can't drift: + never break serialization while an approval prompt is answerable, but never let a + wedged holder park other workers forever (#79719). Resolved once per batch. """ try: from tools.approval import human_wait_ceiling - # human_wait_ceiling is platform-safety-capped (agent/deadline.py - # MAX_SAFE_TIMEOUT_S): a huge approvals.timeout can no longer overflow - # Lock.acquire's time_t on macOS (#83220). Deliberately NOT min()'d - # with _AUTHORIZATION_GATE_LOCK_TIMEOUT_S — the gate must never give - # up while a legitimate approval prompt is still answerable (#79719), - # so a configured approvals.timeout above 360s must extend the gate. + # Safety-capped so a huge approvals.timeout can't overflow Lock.acquire (#83220); + # deliberately NOT min()'d with the fallback so the gate never gives up early (#79719). return human_wait_ceiling() except Exception: return _AUTHORIZATION_GATE_LOCK_TIMEOUT_S @@ -168,8 +145,7 @@ def _authorization_gate_lock_timeout() -> float: class _BatchAbandoned(BaseException): """Raised inside a worker when the batch was abandoned before dispatch. - Derives from BaseException so intermediate ``except Exception`` handlers in - the middleware chain cannot swallow it and dispatch the tool anyway. + BaseException so ``except Exception`` handlers in the middleware chain can't swallow it. """ @@ -193,11 +169,10 @@ def _parse_tool_arguments(raw_arguments: Any) -> tuple[dict, Optional[str]]: def _resolve_concurrent_tool_timeout() -> float | None: - """Resolve the per-batch concurrent tool deadline. + """Resolve the per-batch concurrent tool deadline via the unified resolver (#85125). - Delegates to the unified resolver (#85125): ``timeouts.tools.concurrent_batch`` - in config.yaml wins, the legacy ``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` env var - remains the back-compat bridge, and ``0``/negative still disables the bound. + ``timeouts.tools.concurrent_batch`` wins; ``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` is the + legacy bridge; ``0``/negative disables the bound. """ from agent.deadline import resolve_timeout @@ -214,20 +189,16 @@ def _flush_session_db_after_tool_progress( *, stage: str, ) -> bool: - """Flush tool-call progress before projecting it to any UI surface. + """Flush tool-call progress to the session DB before projecting it to any UI. - Tool execution can perform side effects that terminate or restart the - current Hermes process before the normal turn-end persistence path runs. - Flush the already-appended assistant/tool messages immediately so the - transcript survives destructive-but-valid tool calls. + Tool side effects can kill/restart the process before turn-end persistence runs. """ try: persisted = agent._flush_messages_to_session_db(messages) is not False if not persisted: agent._incremental_persistence_failed = True - # The flush caught its own exception and returned False; the - # classified cause (if any) was captured at the catch site. Only - # fall back to 'unknown' when nothing more specific is recorded. + # Flush recorded any classified cause at the catch site; only default + # to 'unknown' when nothing more specific exists. if getattr(agent, "_last_persistence_error_cause", None) is None: agent._last_persistence_error_cause = "unknown" return persisted @@ -240,11 +211,8 @@ def _flush_session_db_after_tool_progress( def _image_generate_parallel_limit() -> int: - """Return the configured image-generation parallelism cap. - - Image-generation calls are slow enough that concurrent execution is useful, - but backend bursts can hit TTFB or rate-limit failures. Keep the default - intentionally conservative while allowing users to tune it per install. + """Return the configured image-generation parallelism cap (conservative default; + backend bursts hit TTFB/rate-limit failures). """ try: from hermes_cli.config import load_config @@ -286,50 +254,15 @@ def _ra(): def _is_interpreter_shutdown_submit_error(exc: RuntimeError) -> bool: - """Shutdown-race predicate — shared home in ``tools.interpreter_shutdown``. - - Delegates so all sites (cron delivery, conversation-loop retry, tool - submission) recognize both CPython shutdown-message variants instead of - each matching its own substring (the bug class behind #55924/#58720). + """Shutdown-race predicate; delegates to ``tools.interpreter_shutdown`` so every site + recognizes both CPython shutdown-message variants (#55924/#58720). """ from tools.interpreter_shutdown import interpreter_shutting_down return interpreter_shutting_down(exc) -def _emit_terminal_post_tool_call( - agent, - *, - function_name: str, - function_args: dict, - result: Any, - effective_task_id: str, - tool_call_id: str, - duration_ms: int = 0, - status: str | None = None, - error_type: str | None = None, - error_message: str | None = None, - middleware_trace: Optional[list[dict[str, Any]]] = None, -) -> None: - try: - from model_tools import _emit_post_tool_call_hook - _emit_post_tool_call_hook( - function_name=function_name, - function_args=function_args, - result=result, - task_id=effective_task_id or "", - session_id=getattr(agent, "session_id", "") or "", - tool_call_id=tool_call_id or "", - turn_id=getattr(agent, "_current_turn_id", "") or "", - api_request_id=getattr(agent, "_current_api_request_id", "") or "", - duration_ms=duration_ms, - status=status, - error_type=error_type, - error_message=error_message, - middleware_trace=list(middleware_trace or []), - ) - except Exception: - pass +_emit_terminal_post_tool_call = emit_terminal_post_tool_call def _cancelled_tool_result(reason: str = "user interrupt") -> str: @@ -374,17 +307,9 @@ def _emit_cancelled_terminal_post_tool_call( def _tool_search_scoped_names(agent) -> frozenset: """Return the deferrable tool names the session may invoke via tool_call. - The Tool Search unwrap dispatches the underlying tool directly, bypassing - the bridge branch (and its scope check) in - ``model_tools.handle_function_call``. To keep a restricted-toolset session - (subagent, kanban worker, curated gateway session) from reaching tools it - was never granted, the unwrap validates the underlying name against this - set: the deferrable subset of the session's own enabled/disabled toolset - scope. - - Result is cached on the agent and refreshed when the tool registry's - generation changes (e.g. an MCP server reconnects), so the common case is - a dict lookup, not a full tool-defs rebuild on every tool call. + The Tool Search unwrap bypasses the bridge's scope check in + ``model_tools.handle_function_call``, so restricted sessions are validated against + this set. Cached on the agent; refreshed when the registry generation changes. """ try: import model_tools @@ -421,6 +346,60 @@ def _tool_search_scoped_names(agent) -> frozenset: return names +def _canonical_tool_name(function_name: str) -> str: + """Map legacy tool-name aliases (2026-08 renames) BEFORE agent-loop dispatch.""" + from model_tools import _LEGACY_TOOL_ALIASES as _lta + + return _lta.get(function_name, function_name) + + +def _unwrap_tool_search_call( + agent, function_name: str, function_args: dict, *, flatten_probe: bool = False +) -> tuple[str, dict, Optional[str]]: + """Peel the ``tool_call`` bridge so downstream hooks see the underlying tool. + + Checkpointing, guardrails, plugin hooks and the activity feed must observe the real + tool name, not the bridge. ``tool_call.function`` stays untouched for the transcript + and tool_call_id pairing. The unwrap bypasses + handle_function_call's scope check, so session toolset scope is enforced HERE. + Returns ``(name, args, scope_block)``; ``scope_block`` is the block message when + the underlying tool is out of scope or its args fail the deferred-schema probe + (``flatten_probe`` collapses the probe's JSON payload to one plain string for + callers that wrap the message in ``{"error": ...}``). + """ + scope_block: Optional[str] = None + try: + from tools import tool_search as _ts + if function_name == _ts.TOOL_CALL_NAME: + underlying, underlying_args, err = _ts.resolve_underlying_call(function_args) + if not err and underlying: + if underlying in _tool_search_scoped_names(agent): + # Validate before unwrapping: the generic bridge hides the concrete + # parameter schema from provider-native tool-call validation. + probe_err = _ts.validate_deferred_call_args(underlying, underlying_args) + if probe_err is None: + return underlying, underlying_args, None + scope_block = probe_err + if flatten_probe: + try: + probe = json.loads(probe_err) + scope_block = ( + f"{probe.get('error', '')} Parameters schema: " + f"{json.dumps(probe.get('parameters', {}), ensure_ascii=False)}. " + f"{probe.get('hint', '')}" + ).strip() + except Exception: + scope_block = probe_err + else: + scope_block = ( + f"'{underlying}' is not available in this session. " + "Use tool_search to find tools you can call." + ) + except Exception: + pass + return function_name, function_args, scope_block + + @dataclass class _ManagedToolResult: result: Any @@ -437,33 +416,20 @@ class _ToolTimeoutResult(str): class _ToolCancelledResult(str): """Marker for a synthesized sequential-tool user-interrupt result. - Like ``_ToolTimeoutResult``, the executor already emitted the terminal - post_tool_call event for this call (status="cancelled"), so downstream - emission must be suppressed — an abandoned worker finishing late must not - report success for a call the user already cancelled. + The terminal post_tool_call event was already emitted (status=cancelled), so a + late-finishing abandoned worker must not report success. """ class _ConcurrentToolAuthorizationGate: """Serialize policy prompts and exclude human approval waits from batch deadlines. - Serialization keeps concurrent approval prompts from interleaving on the - user's screen. The acquire is BOUNDED: a worker wedged inside the gate (a - hanging ``pre_tool_call`` plugin, or an approval round-trip to a client - that went away) must not park every other worker forever. On expiry the - worker runs its prompt unserialized — worst case is interleaved prompts, - strictly better than permanent starvation (same tradeoff as the - start-order gate, #79705). + The acquire is BOUNDED: on expiry the worker prompts unserialized rather than + starving the batch behind a wedged plugin/approval client (#79705). Deadline exclusion is measured at the SOURCE of the human wait - (``tools.approval.human_wait_seconds``: the CLI prompt and the gateway - approval poll loop mark their own blocking windows), NOT as residency in - this gate. Gate residency is arbitrary code — using it as the exclusion - signal let a wedged plugin grow the exclusion 1:1 with wall clock, keeping - the batch deadline's ``remaining`` constant so it never fired and the turn - hung forever (#79719). A wedged plugin now contributes nothing to the - exclusion and the batch times out normally, while a genuine approval wait - (which can legitimately exceed any fixed bound) is still excluded in full. + (``tools.approval.human_wait_seconds``), NOT as gate residency: residency-based + exclusion let a wedged plugin keep the deadline from ever firing (#79719). """ def __init__( @@ -483,9 +449,8 @@ class _ConcurrentToolAuthorizationGate: try: from tools.approval import get_current_session_key - # Snapshot the batch's session identity on the SUBMITTING - # thread: excluded_seconds() is polled from the batch wait - # loop, whose context may differ from the workers'. + # Snapshot on the SUBMITTING thread: excluded_seconds() is polled + # from the batch wait loop, whose context may differ from workers'. self._session_key = get_current_session_key() except Exception: logger.debug( @@ -535,9 +500,8 @@ def _managed_values( ) -# Cadence for the in-flight tool activity heartbeat. Must stay far below the -# gateway turn-inactivity timeout (default 1800s) so a silent-but-healthy -# tool call never looks idle to the watchdog. +# Heartbeat cadence; must stay far below the gateway turn-inactivity timeout +# (default 1800s) so a silent-but-healthy tool never looks idle. _TOOL_ACTIVITY_HEARTBEAT_INTERVAL_S = 30.0 @@ -547,28 +511,11 @@ def _run_tool_activity_heartbeat( label: str, interval: float = _TOOL_ACTIVITY_HEARTBEAT_INTERVAL_S, ) -> None: - """Refresh the agent's activity clock while a tool call is in flight. + """Daemon thread that stamps ``agent._touch_activity`` every ``interval`` seconds + until ``stop_event`` is set. - The gateway's turn-inactivity watchdog - (``gateway/run.py::_watch_gateway_turn_inactivity``) abandons a turn - once ``seconds_since_activity`` exceeds the inactivity timeout - (default 30 min). Activity is stamped when a tool *starts* and when it - *completes*, but a tool call that runs silently for 30+ minutes - (quiet builds, long pytest suites, large downloads, network waits that - emit no output) previously froze the clock at "executing tool: " - and the watchdog hard-abandoned a turn that was still making progress, - reaping the tool's processes mid-execution. - - This daemon thread touches ``agent._touch_activity`` every ``interval`` - seconds until ``stop_event`` is set (the tool call returned), so the - gateway keeps seeing a live turn for the whole duration of the call. - - A tool that truly hangs is still bounded by the tool layer's own - timeouts (terminal ``timeout`` default 180s, the concurrent batch - deadline ~420s), so the heartbeat only extends the turn's life for as - long as the tool call is legitimately executing — it does not unbind - wedged tools. The 30-min gateway backstop remains for turns whose - agent loop itself stalls (no API call, no tool call in flight). + Keeps the gateway turn-inactivity watchdog (default 30 min) from abandoning a turn + whose tool runs silently. Wedged tools stay bounded by the tool layer's own timeouts. """ try: @@ -649,12 +596,7 @@ def _run_agent_tool_execution_middleware( block_msg, modified_args = _dispatch_pre_tool_call_hooks( function_name, final_args, - task_id=effective_task_id or "", - session_id=getattr(agent, "session_id", "") or "", - tool_call_id=tool_call_id or "", - turn_id=getattr(agent, "_current_turn_id", "") or "", - api_request_id=getattr(agent, "_current_api_request_id", "") - or "", + **tool_hook_ids(agent, effective_task_id, tool_call_id), middleware_trace=list(state["middleware_trace"]), ) if modified_args is not None: @@ -713,12 +655,8 @@ def _run_agent_tool_execution_middleware( _advance_start_order(_begin) - # Keep the gateway turn-inactivity watchdog from abandoning a turn - # whose tool call runs silently for longer than the inactivity - # timeout (#84491): stamp activity periodically while the tool is - # in flight, not just at start/completion. Both the sequential and - # the concurrent paths funnel through here, so a single heartbeat - # covers every tool. + # Heartbeat while the tool is in flight so the gateway inactivity watchdog + # doesn't abandon a silent-but-live turn (#84491); covers both executor paths. _hb_stop = threading.Event() _hb_thread = threading.Thread( target=_run_tool_activity_heartbeat, @@ -739,11 +677,7 @@ def _run_agent_tool_execution_middleware( function_name, relay_args, skip_relay=True, - task_id=effective_task_id or "", - session_id=getattr(agent, "session_id", "") or "", - tool_call_id=tool_call_id or "", - turn_id=getattr(agent, "_current_turn_id", "") or "", - api_request_id=getattr(agent, "_current_api_request_id", "") or "", + **tool_hook_ids(agent, effective_task_id, tool_call_id), ) request_args = ( request_result.payload @@ -759,11 +693,7 @@ def _run_agent_tool_execution_middleware( next_args if isinstance(next_args, dict) else request_args ), original_args=function_args, - task_id=effective_task_id or "", - session_id=getattr(agent, "session_id", "") or "", - tool_call_id=tool_call_id or "", - turn_id=getattr(agent, "_current_turn_id", "") or "", - api_request_id=getattr(agent, "_current_api_request_id", "") or "", + **tool_hook_ids(agent, effective_task_id, tool_call_id), ) result, _relay_args = relay_tools.execute( @@ -787,26 +717,20 @@ def _run_agent_tool_execution_middleware( ) -# How often the sequential-tool wait loop wakes to check for a user -# interrupt while the worker runs. Short enough that /stop or a redirect -# lands within ~1s even when the tool itself never polls is_interrupted(). +# Sequential wait-loop interrupt poll cadence: /stop lands within ~1s even when +# the tool never polls is_interrupted(). _SEQUENTIAL_INTERRUPT_POLL_SECONDS = 1.0 def _resolve_sequential_tool_timeout() -> float | None: """Deadline for one sequential tool call (#85125 Phase 2a). - ``timeouts.tools.sequential_call`` in config.yaml wins; when unset, the - sequential path inherits the concurrent batch deadline (same value, same - ``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` legacy bridge) so the two executor - paths cannot drift apart by default. ``0``/negative disables the bound. + ``timeouts.tools.sequential_call`` wins; unset inherits the concurrent batch deadline + so the two paths can't drift. ``0``/negative disables the bound. - NOTE: this path deliberately does NOT use ``agent.deadline.run_bounded_sync``. - The sequential/concurrent executors extend their deadline dynamically while - a human approval prompt is open (``_ConcurrentToolAuthorizationGate`` - excluded seconds — a MUST-preserve invariant) and touch agent activity - mid-wait; the shared primitive is fixed-deadline by design. Simpler call - sites migrate onto the primitive; these two stay symmetric with each other. + Deliberately NOT ``agent.deadline.run_bounded_sync``: both executors extend their + deadline while an approval prompt is open (MUST-preserve), which the fixed-deadline + primitive can't express. """ from agent.deadline import resolve_timeout @@ -830,10 +754,8 @@ def _run_sequential_tool_execution_middleware( ) -> _ManagedToolResult: """Run one sequential call with the concurrent executor's deadline. - Interactive input tools such as ``clarify`` wait on a human. Their own - timeout (``agent.clarify_timeout``: default 3600s, or unlimited when - ``<= 0``) owns that wait. Applying the generic tool deadline here would - return ``tool_timeout`` while the prompt and worker stay active. + Interactive tools (``clarify``) own their wait via ``agent.clarify_timeout``; the + generic deadline would report ``tool_timeout`` while the prompt is still live. """ timeout_s = _resolve_sequential_tool_timeout() kwargs = { @@ -873,10 +795,8 @@ def _run_sequential_tool_execution_middleware( executor = DaemonThreadPoolExecutor(max_workers=1) future = executor.submit(propagate_context_to_thread(_run)) - # ``timeout_s`` disabled (None) still runs on the worker: the wait loop - # below is what makes a non-cooperative tool interruptible at all, so - # "no deadline" must not mean "no interrupt checks" (#86xxx class fix — - # sequential path previously blocked until the tool returned). + # Disabled timeout still runs on the worker: this wait loop is what makes a + # non-cooperative tool interruptible, so no deadline must not mean no interrupt checks. deadline = time.monotonic() + timeout_s if timeout_s is not None else None started = time.monotonic() timed_out = False @@ -918,9 +838,8 @@ def _run_sequential_tool_execution_middleware( ) except Exception: pass - # Give a cooperative tool a moment to notice its per-thread - # interrupt bit and return a real result (mirrors the concurrent - # path's 3s grace). + # Grace for a cooperative tool to notice its interrupt bit (mirrors the + # concurrent path's 3s). concurrent.futures.wait([future], timeout=3.0) if future.done() and not future.cancelled(): return future.result() @@ -1095,21 +1014,115 @@ def _begin_tool_execution( pass +def _append_finalized_tool_result( + agent, + messages: list, + *, + function_name: str, + function_args: dict, + function_result, + tool_call_id: str, + effective_task_id: str, + budget: BudgetConfig, + effect_disposition=None, +): + """Persist/spill, hint, wrap and append one tool result; flush the session DB. + + Returns ``(function_result, tool_message, risk_metadata)`` — ``function_result`` is the + persisted/hinted content — or ``None`` when the incremental flush failed (the caller + must stop the batch). + """ + if not _is_multimodal_tool_result(function_result): + function_result = maybe_persist_tool_result( + content=function_result, + tool_name=function_name, + tool_use_id=tool_call_id, + env=get_active_env(effective_task_id), + config=budget, + ) + _record_persisted_path_for_stub(agent, tool_call_id, function_result) + + subdir_hints = agent._subdirectory_hints.check_tool_call(function_name, function_args) + if subdir_hints: + if _is_multimodal_tool_result(function_result): + # Append the hint to the text summary part so the model still sees it; + # don't touch the image blocks. + _append_subdir_hint_to_multimodal(function_result, subdir_hints) + else: + function_result += subdir_hints + + # Unwrap _multimodal dicts to an OpenAI-style content list; text-only servers + # get a string-safe fallback so a rejected image result never poisons history. + _tool_content = agent._tool_result_content_for_active_model(function_name, function_result) + tool_message = make_tool_result_message( + function_name, + _tool_content, + tool_call_id, + effect_disposition=effect_disposition, + ) + messages.append(tool_message) + if not _flush_session_db_after_tool_progress( + agent, + messages, + stage=f"tool result {function_name}", + ): + return None + return function_result, tool_message, tool_message.get("_tool_output_risk") + + +def _emit_tool_completed_progress(agent, function_name: str, *, duration: float, is_error: bool, result) -> None: + """``tool.completed`` UI projection; downstream of the canonical append so resume + can reconstruct the result even if the UI bridge dies mid-projection.""" + if not agent.tool_progress_callback: + return + try: + agent.tool_progress_callback( + "tool.completed", function_name, None, None, + duration=duration, is_error=is_error, result=result, + ) + except Exception as cb_err: + logging.debug("Tool progress callback error: %s", cb_err) + + +def _emit_tool_complete_and_risk( + agent, *, function_name: str, function_args: dict, tool_call_id: str, result, risk_metadata, blocked: bool +) -> None: + """Fire ``tool_complete_callback`` (unless blocked) then the ``tool.output_risk`` projection.""" + if not blocked and agent.tool_complete_callback: + try: + display_args = _redact_tool_args_for_display(function_name, function_args) or function_args + agent.tool_complete_callback(tool_call_id, function_name, display_args, result) + except Exception as cb_err: + logging.debug("Tool complete callback error: %s", cb_err) + + if ( + risk_metadata is not None + and risk_metadata.get("risk") != "low" + and agent.tool_progress_callback + ): + try: + agent.tool_progress_callback( + "tool.output_risk", + function_name, + None, + None, + tool_call_id=tool_call_id, + risk_metadata=risk_metadata, + ) + except Exception as cb_err: + logging.debug("Tool output risk callback error: %s", cb_err) + + def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effective_task_id: str, api_call_count: int = 0, *, finalize: bool = True) -> None: - """Execute multiple tool calls concurrently using a thread pool. + """Execute tool calls concurrently; results are appended in original call order. - Results are collected in the original tool-call order and appended to - messages so the API sees them in the expected sequence. - - ``finalize=False`` skips the end-of-batch aggregate budget enforcement - and /steer injection — used when this call is one segment of a larger - mixed batch and the segmented dispatcher owns the turn-end work. + ``finalize=False`` skips end-of-batch budget enforcement and /steer injection (the + segmented dispatcher owns turn-end work). """ tool_calls = assistant_message.tool_calls num_tools = len(tool_calls) - # Resolve the context-scaled tool-output budget once per turn (cheap, but - # avoids rebuilding it per result inside the loop below). + # Resolve the context-scaled tool-output budget once per turn, not per result. _tool_budget = _budget_for_agent(agent) # ── Pre-flight: interrupt check ────────────────────────────────── @@ -1145,17 +1158,11 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe ) return - # ── Parse args + pre-execution bookkeeping ─────────────────────── - # (tool call, resolved name, parsed args, middleware trace, parse error, - # tool-search scope block) + # ── Parse args + pre-execution bookkeeping ──────────────────────────── + # (tool call, name, args, middleware trace, parse error, tool-search scope block) parsed_calls = [] for tool_call in tool_calls: - function_name = tool_call.function.name - # Legacy tool-name aliases (2026-08 renames) — map BEFORE the - # agent-loop branches (todo_list etc. dispatch above the registry). - from model_tools import _LEGACY_TOOL_ALIASES as _lta - function_name = _lta.get(function_name, function_name) - + function_name = _canonical_tool_name(tool_call.function.name) function_args, malformed_args_result = _parse_tool_arguments( tool_call.function.arguments ) @@ -1173,45 +1180,9 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe ) continue - # ── Tool Search unwrap ──────────────────────────────────────── - # When the model invokes the tool_call bridge, peel it open so - # every downstream check (checkpointing, guardrails, plugin - # pre-tool-call hooks, the display/activity feed, the post-call - # callback) sees the underlying tool — not the bridge. This is - # the OpenClaw lesson: hooks must observe the real tool name. - # - # The original tool_call entry on ``tool_call.function`` is left - # untouched so the conversation transcript and the matching - # tool_call_id are preserved exactly as the model emitted them. - # - # Scope gate: the unwrap dispatches the underlying tool directly - # (bypassing the bridge branch in handle_function_call and its - # scope check), so we enforce session toolset scope HERE. A tool - # the session was not granted is rejected before any checkpoint, - # hook, or dispatch fires. - _ts_scope_block = None - try: - from tools import tool_search as _ts - if function_name == _ts.TOOL_CALL_NAME: - _underlying, _underlying_args, _err = _ts.resolve_underlying_call(function_args) - if not _err and _underlying: - if _underlying in _tool_search_scoped_names(agent): - # Validate before unwrapping: the generic bridge hides - # the concrete parameter schema from provider-native - # tool-call validation. - _probe_err = _ts.validate_deferred_call_args(_underlying, _underlying_args) - if _probe_err is not None: - _ts_scope_block = _probe_err - else: - function_name = _underlying - function_args = _underlying_args - else: - _ts_scope_block = ( - f"'{_underlying}' is not available in this session. " - "Use tool_search to find tools you can call." - ) - except Exception: - pass + function_name, function_args, _ts_scope_block = _unwrap_tool_search_call( + agent, function_name, function_args + ) parsed_calls.append( (tool_call, function_name, function_args, [], None, _ts_scope_block) @@ -1231,9 +1202,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe start_condition = threading.Condition() next_start_order = 0 - # Set once the batch is abandoned (deadline or interrupt) so a worker parked - # at the start-order gate exits immediately instead of waking up minutes - # later and dispatching a tool the turn has already reported as timed out. + # Set once the batch is abandoned so gate-parked workers exit instead of + # dispatching a tool the turn already reported as timed out. batch_abandoned = threading.Event() authorization_gate = _ConcurrentToolAuthorizationGate() @@ -1243,10 +1213,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe with start_condition: start_condition.notify_all() - # The gate bound must sit UNDER the batch deadline, otherwise the deadline - # fires first and the parked workers are still falsely reported as timed - # out without ever starting — the very bug this gate timeout fixes. A - # disabled deadline (None) keeps the stock bound rather than waiting forever. + # The gate bound must sit UNDER the batch deadline, else parked workers are falsely + # reported timed out without starting. A disabled deadline keeps the stock bound. def _start_order_gate_timeout(batch_timeout: float | None) -> float: if batch_timeout is None: return _START_ORDER_GATE_TIMEOUT_S @@ -1258,21 +1226,9 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe """Serialize dispatch by submit order. Returns False if abandoned.""" nonlocal next_start_order with start_condition: - # Bounded wait: a tool that wedges during its dispatch must not - # park every later-ordered worker forever. Without the timeout, - # one blocking dispatch starves the whole batch (the parked tools - # then get falsely reported as "timed out" by the batch deadline - # despite never having started) and the parked threads leak - # permanently after the batch is abandoned — f.cancel() cannot - # cancel running threads and nothing ever notifies the condition - # again. On expiry, proceed out of order: the worst case is - # interleaved approval prompts, strictly better than permanent - # starvation. The >= predicate (rather than ==) lets one worker's - # timeout-jump release every skipped worker immediately instead - # of each burning its own full timeout; max() keeps the counter - # monotonic when workers advance out of order. batch_abandoned - # short-circuits the wait so an abandoned batch releases its - # parked workers in milliseconds instead of one gate timeout. + # Bounded wait so one wedged dispatch can't starve/leak later-ordered workers; + # on expiry proceed out of order (interleaved prompts beat starvation). + # >= (not ==) releases every skipped worker at once; batch_abandoned short-circuits. in_order = start_condition.wait_for( lambda: next_start_order >= order or batch_abandoned.is_set(), timeout=( @@ -1319,16 +1275,13 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe start_order, ): """Worker function executed in a thread.""" - # Register this worker tid so the agent can fan out an interrupt - # to it — see AIAgent.interrupt(). Must happen first thing, and - # must be paired with discard + clear in the finally block. + # Register this worker tid for interrupt fan-out (AIAgent.interrupt()); must be + # first and paired with discard + clear in finally. _worker_tid = threading.current_thread().ident with agent._tool_worker_threads_lock: agent._tool_worker_threads.add(_worker_tid) - # Race: if the agent was interrupted between fan-out (which - # snapshotted an empty/earlier set) and our registration, apply - # the interrupt to our own tid now so is_interrupted() inside - # the tool returns True on the next poll. + # Race: interrupt may have fanned out before our registration; apply it + # to our own tid now. if agent._interrupt_requested: try: _ra()._set_interrupt( @@ -1338,18 +1291,15 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe ) except Exception: pass - # Set the activity callback on THIS worker thread so - # _wait_for_process (terminal commands) can fire heartbeats. - # The callback is thread-local; the main thread's callback - # is invisible to worker threads. + # Activity callback is thread-local; set it on THIS worker so + # _wait_for_process heartbeats fire. try: from tools.environments.base import set_activity_callback set_activity_callback(agent._touch_activity) except Exception: pass - # Approval/sudo callbacks (thread-local) and the agent turn's - # ContextVars are propagated by propagate_context_to_thread() at the - # submit site below (GHSA-qg5c-hvr5-hjgr, #13617). + # Approval/sudo callbacks and turn ContextVars are propagated by + # propagate_context_to_thread() at submit (GHSA-qg5c-hvr5-hjgr, #13617). start = time.time() tool_call_id = _pairing_tool_call_id(tool_call) blocked = False @@ -1408,11 +1358,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe blocked = managed.blocked dispatched = managed.dispatched except _BatchAbandoned: - # The batch was abandoned while we were parked at the start-order - # gate. The main thread already synthesized this tool's result - # (timeout/cancelled) and moved on, so write nothing: a late - # results[index] write, post_tool_call emit, or progress print - # would double-report a tool_call_id the turn already closed. + # Abandoned at the start-order gate: the main thread already synthesized + # this result, so write/emit nothing (would double-report the tool_call_id). logger.info( "tool %s abandoned at start-order gate; skipping dispatch", function_name, @@ -1474,19 +1421,14 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe middleware_trace, ) finally: - # Teardown advance: keep the counter moving for any later-ordered - # worker. Never let the abandonment signal escape from here — the - # worker is already unwinding and the turn owns the result. + # Teardown advance keeps later-ordered workers moving; never let the + # abandonment signal escape here. try: _advance_start() except _BatchAbandoned: pass - # Tear down worker-tid tracking. Clear any interrupt bit we may - # have set so the next task scheduled onto this recycled tid - # starts with a clean slate. This MUST be in a finally block - # because BaseException subclasses (CancelledError, KeyboardInterrupt) - # bypass ``except Exception`` and would otherwise leak the tid - # into _interrupted_threads, poisoning the recycled thread. + # Tear down tid tracking and clear any interrupt bit so a recycled tid starts + # clean. MUST be in finally: BaseException subclasses bypass ``except Exception``. with agent._tool_worker_threads_lock: agent._tool_worker_threads.discard(_worker_tid) try: @@ -1515,11 +1457,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe deadline = time.monotonic() + timeout_s if timeout_s is not None else None if runnable_calls: max_workers = _max_workers_for_tool_batch(runnable_calls) - # Daemon workers: an interrupted/timed-out batch is abandoned with - # shutdown(wait=False), but stdlib ThreadPoolExecutor workers are - # non-daemon and registered in concurrent.futures' atexit hook, - # which joins them unconditionally — so one wedged tool thread - # would block interpreter exit forever (multi-minute CLI exits). + # Daemon workers: stdlib ThreadPoolExecutor's atexit join would let one + # wedged tool thread block interpreter exit forever. from tools.daemon_pool import DaemonThreadPoolExecutor executor = DaemonThreadPoolExecutor(max_workers=max_workers) abandon_executor = False @@ -1527,9 +1466,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe for submit_index, (i, tc, name, args, scope_block) in enumerate( runnable_calls ): - # Propagate the agent turn's ContextVars (e.g. - # _approval_session_key) AND thread-local approval/sudo - # callbacks into the worker thread; clears callbacks on exit. + # Propagate turn ContextVars and thread-local approval/sudo + # callbacks into the worker; clears callbacks on exit. try: f = executor.submit( propagate_context_to_thread(_run_tool), @@ -1576,11 +1514,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe futures.append(f) future_to_index[f] = i - # Wait for all to complete with periodic heartbeats so the - # gateway's inactivity monitor doesn't kill us during long - # concurrent tool batches. Also check for user interrupts - # so we don't block indefinitely when the user sends /stop - # or a new message during concurrent tool execution. + # Wait with periodic heartbeats (gateway inactivity monitor) and + # interrupt checks (/stop or a new message). _conc_start = time.time() _interrupt_logged = False while True: @@ -1630,9 +1565,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe ) for f in not_done: f.cancel() - # Release gate-parked workers before the interrupt - # fan-out so none of them wakes up later and dispatches - # a tool this loop just reported as timed out. + # Release gate-parked workers before interrupt fan-out so none + # later dispatches a tool just reported as timed out. _abandon_batch() with agent._tool_worker_threads_lock: worker_tids = list(agent._tool_worker_threads) @@ -1643,11 +1577,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe pass break - # Check for interrupt — the per-thread interrupt signal - # already causes individual tools (terminal, execute_code) - # to abort, but tools without interrupt checks (web_search, - # read_file) will run to completion. Cancel any futures - # that haven't started yet so we don't block on them. + # Tools without interrupt checks (web_search, read_file) run to + # completion; cancel unstarted futures so we don't block on them. if agent._interrupt_requested: abandon_executor = True if not _interrupt_logged: @@ -1680,23 +1611,18 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe f"{len(not_done)} remaining: {', '.join(_still_running[:3])})" ) finally: - # Belt-and-braces: any exit from the wait loop that abandoned - # the batch must release gate-parked workers, including the - # exception path that never reaches the branches above. + # Any abandoning exit from the wait loop (including the exception + # path) must release gate-parked workers. if abandon_executor: _abandon_batch() - # On abandon (interrupt or deadline) we intentionally do NOT - # join hung workers: wait=False returns immediately and - # cancel_futures drops queued-but-unstarted work. A wedged tool - # thread is left running detached — the deliberate tradeoff vs. - # deadlocking the whole batch. Normal completion joins (wait=True). + # On abandon do NOT join hung workers: a wedged thread is left detached + # rather than deadlocking the batch. Normal completion joins. executor.shutdown( wait=not abandon_executor, cancel_futures=abandon_executor, ) finally: if spinner: - # Build a summary message for the spinner stop completed = sum(1 for r in results if r is not None) total_dur = sum(r[3] for r in results if r is not None) spinner.stop(f"⚡ {completed}/{num_tools} tools completed in {total_dur:.1f}s total") @@ -1710,10 +1636,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe blocked = False is_error = True progress_function_name = name - # A worker can finish and write results[i] in the window between the - # deadline snapshot (timed_out_indices, taken from not_done) and this - # loop. Prefer that real result over a fabricated timeout message — the - # tool genuinely succeeded, just slightly late. + # A worker may finish between the deadline snapshot and this loop; + # prefer its real result over a fabricated timeout. effect_disposition = None if i in timed_out_indices and r is None: suffix = f"{timeout_s:.1f}s" if timeout_s is not None else "the configured timeout" @@ -1799,9 +1723,8 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe result_preview = _err_text[:200] if len(_err_text) > 200 else _err_text logger.warning("Tool %s returned error (%.2fs): %s", function_name, tool_duration, result_preview) - # Track file-mutation outcome for the turn-end verifier. - # `blocked` calls never actually ran — don't let a guardrail - # block count as either a failure or a success. + # Track file-mutation outcome for the turn-end verifier; blocked calls + # never ran, so they count as neither failure nor success. if not blocked: try: agent._record_file_mutation_result( @@ -1819,62 +1742,27 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe agent._touch_activity(f"tool completed: {name} ({tool_duration:.1f}s){_status_suffix}") display_function_result = function_result - function_result = maybe_persist_tool_result( - content=function_result, - tool_name=name, - tool_use_id=tool_call_id, - env=get_active_env(effective_task_id), - config=_tool_budget, - ) if not _is_multimodal_tool_result(function_result) else function_result - _record_persisted_path_for_stub(agent, tool_call_id, function_result) - - subdir_hints = agent._subdirectory_hints.check_tool_call(name, args) - if subdir_hints: - if _is_multimodal_tool_result(function_result): - # Append the hint to the text summary part so the model - # still sees it; don't touch the image blocks. - _append_subdir_hint_to_multimodal(function_result, subdir_hints) - else: - function_result += subdir_hints - - # Unwrap _multimodal dicts to an OpenAI-style content list so any - # vision-capable provider receives [{type:text},{type:image_url}] - # rather than a raw Python dict. The Anthropic adapter already - # accepts content lists; vision-capable OpenAI-compatible servers - # (mlx-vlm, GPT-4o, …) accept image_url in tool messages natively. - # Text-only servers get a string-safe fallback here so a rejected - # image tool result never poisons canonical session history. - # String results pass through unchanged. - _tool_content = agent._tool_result_content_for_active_model(name, function_result) - tool_message = make_tool_result_message( - name, - _tool_content, - tool_call_id, - effect_disposition=effect_disposition, - ) - messages.append(tool_message) - risk_metadata = tool_message.get("_tool_output_risk") - if not _flush_session_db_after_tool_progress( + finalized = _append_finalized_tool_result( agent, messages, - stage=f"tool result {name}", - ): + function_name=name, + function_args=args, + function_result=function_result, + tool_call_id=tool_call_id, + effective_task_id=effective_task_id, + budget=_tool_budget, + effect_disposition=effect_disposition, + ) + if finalized is None: return + function_result, _tool_message, risk_metadata = finalized - # Every completion surface is downstream of the canonical append. If - # the UI bridge or process dies while projecting one of these events, - # resume can reconstruct the tool result that was already visible. - if not blocked and agent.tool_progress_callback: - try: - agent.tool_progress_callback( - "tool.completed", progress_function_name, None, None, - duration=tool_duration, is_error=is_error, - result=display_function_result, - ) - except Exception as cb_err: - logging.debug("Tool progress callback error: %s", cb_err) + if not blocked: + _emit_tool_completed_progress( + agent, progress_function_name, + duration=tool_duration, is_error=is_error, result=display_function_result, + ) - # Print cute message per tool if agent._should_emit_quiet_tool_messages(): cute_msg = _get_cute_tool_message_impl( name, args, tool_duration, result=display_function_result, @@ -1889,59 +1777,34 @@ def execute_tool_calls_concurrent(agent, assistant_message, messages: list, effe response_preview = _preview_str[:agent.log_prefix_chars] + "..." if len(_preview_str) > agent.log_prefix_chars else _preview_str print(f" ✅ Tool {i+1} completed in {tool_duration:.2f}s - {response_preview}") - if not blocked and agent.tool_complete_callback: - try: - display_args = _redact_tool_args_for_display(name, args) or args - agent.tool_complete_callback( - tool_call_id, name, display_args, display_function_result, - ) - except Exception as cb_err: - logging.debug("Tool complete callback error: %s", cb_err) + _emit_tool_complete_and_risk( + agent, + function_name=name, + function_args=args, + tool_call_id=tool_call_id, + result=display_function_result, + risk_metadata=risk_metadata, + blocked=blocked, + ) - if ( - risk_metadata is not None - and risk_metadata.get("risk") != "low" - and agent.tool_progress_callback - ): - try: - agent.tool_progress_callback( - "tool.output_risk", - name, - None, - None, - tool_call_id=tool_call_id, - risk_metadata=risk_metadata, - ) - except Exception as cb_err: - logging.debug("Tool output risk callback error: %s", cb_err) - - # ── Per-turn aggregate budget enforcement ───────────────────────── - # Keep /steer pending until the final post-budget drain below. The model - # cannot observe a partial batch, while an early drain can be discarded - # when aggregate budget enforcement replaces that tool result. + # ── Per-turn aggregate budget enforcement ────────────────────────── + # Keep /steer pending until the post-budget drain: an early drain could be + # discarded when budget enforcement replaces that tool result. num_tools = len(parsed_calls) if finalize and num_tools > 0: turn_tool_msgs = messages[-num_tools:] enforce_turn_budget(turn_tool_msgs, env=get_active_env(effective_task_id), config=_tool_budget) - # ── /steer injection ────────────────────────────────────────────── - # Append any pending user steer text to the last tool result so the - # agent sees it on its next iteration. Runs AFTER budget enforcement - # so the steer marker is never truncated. See steer() for details. + # ── /steer injection ──────────────────────────────────────────────── + # AFTER budget enforcement so the steer marker is never truncated; see steer(). if finalize and num_tools > 0: agent._apply_pending_steer_to_tool_results(messages, num_tools) def _append_cancelled_tool_results(messages: list, tool_calls, *, reason: str) -> None: - """Append a cancelled ``tool`` result for each call in ``tool_calls``. - - Used when a hard interrupt (KeyboardInterrupt / BaseException) aborts the - sequential executor mid-batch. Without this, the loop re-raises leaving the - assistant tool-call turn with no matching tool results — a message-role - alternation violation that malforms the next provider request. Mirrors the - cooperative-interrupt skip block and the concurrent path, both of which - already emit a result for every call_id. + """Append a cancelled ``tool`` result for each call so a hard interrupt never leaves + the assistant tool-call turn without matching results (role-alternation violation). """ for tc in tool_calls: name = getattr(getattr(tc, "function", None), "name", "") or "tool" @@ -1953,18 +1816,44 @@ def _append_cancelled_tool_results(messages: list, tool_calls, *, reason: str) - )) -def execute_tool_calls_sequential(agent, assistant_message, messages: list, effective_task_id: str, api_call_count: int = 0, *, finalize: bool = True) -> None: - """Execute tool calls sequentially (original behavior). Used for single calls or interactive tools. +def _start_quiet_tool_spinner(agent, function_name: str, function_args: dict, *, gate: bool = True): + """Start the quiet-mode kawaii spinner for one tool call, or return None. - ``finalize=False`` skips the end-of-batch aggregate budget enforcement - and /steer injection — used when this call is one segment of a larger - mixed batch and the segmented dispatcher owns the turn-end work. + ``gate=False`` skips ``_should_start_quiet_spinner`` (context-engine tools always spin). """ - # Resolve the context-scaled tool-output budget once per turn. + if not agent._should_emit_quiet_tool_messages(): + return None + if gate and not agent._should_start_quiet_spinner(): + return None + face = random.choice(KawaiiSpinner.get_waiting_faces()) + emoji = _get_tool_emoji(function_name) + display_args = _redact_tool_args_for_display(function_name, function_args) or function_args + preview = _build_tool_label(function_name, display_args) or function_name + spinner = KawaiiSpinner(f"{face} {emoji} {preview}", spinner_type='dots', print_fn=agent._print_fn) + spinner.start() + return spinner + + +def _finish_quiet_tool_spinner(agent, spinner, function_name: str, function_args: dict, tool_duration: float, result) -> None: + """Stop the spinner with the cute completion line, or print it when no spinner ran.""" + cute_msg = _get_cute_tool_message_impl(function_name, function_args, tool_duration, result=result) + if spinner: + spinner.stop(cute_msg) + elif agent._should_emit_quiet_tool_messages(): + agent._vprint(f" {cute_msg}") + + +def execute_tool_calls_sequential(agent, assistant_message, messages: list, effective_task_id: str, api_call_count: int = 0, *, finalize: bool = True) -> None: + """Execute tool calls sequentially (single calls or interactive tools). + + ``finalize=False`` skips end-of-batch budget enforcement and /steer injection (the + segmented dispatcher owns turn-end work). + """ + # Resolve the context-scaled tool-output budget once per turn, not per result. _tool_budget = _budget_for_agent(agent) - # Keep every runtime-tool branch on one bounded execution funnel without - # duplicating timeout policy across the branch-specific callbacks below. + # One bounded execution funnel for every runtime-tool branch; no duplicated + # timeout policy in the callbacks below. def _run_agent_tool_execution_middleware(agent, **kwargs): return _run_sequential_tool_execution_middleware(agent, **kwargs) @@ -1972,9 +1861,8 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe tool_call_id = _pairing_tool_call_id(tool_call) if getattr(agent, "_incremental_persistence_failed", False): return - # SAFETY: check interrupt BEFORE starting each tool. - # If the user sent "stop" during a previous tool's execution, - # do NOT start any more tools -- skip them all immediately. + # SAFETY: check interrupt BEFORE each tool so a "stop" during the previous + # tool skips all remaining ones. if agent._interrupt_requested: remaining_calls = assistant_message.tool_calls[i-1:] if remaining_calls: @@ -2010,12 +1898,7 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe return break - function_name = tool_call.function.name - # Legacy tool-name aliases (2026-08 renames) — map BEFORE the - # agent-loop branches (todo_list etc. dispatch above the registry). - from model_tools import _LEGACY_TOOL_ALIASES as _lta - function_name = _lta.get(function_name, function_name) - + function_name = _canonical_tool_name(tool_call.function.name) function_args, malformed_args_result = _parse_tool_arguments( tool_call.function.arguments ) @@ -2046,42 +1929,9 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe return continue - # Tool Search unwrap — see execute_tool_calls_concurrent for full - # rationale, including the scope gate (the unwrap dispatches the - # underlying tool directly, so session toolset scope is enforced here). - _ts_scope_block: Optional[str] = None - try: - from tools import tool_search as _ts - if function_name == _ts.TOOL_CALL_NAME: - _underlying, _underlying_args, _err = _ts.resolve_underlying_call(function_args) - if not _err and _underlying: - if _underlying in _tool_search_scoped_names(agent): - # Validate before unwrapping: the generic bridge hides - # the concrete parameter schema from provider-native - # tool-call validation. - _probe_err = _ts.validate_deferred_call_args(_underlying, _underlying_args) - if _probe_err is not None: - # This path wraps _block_msg in {"error": ...} — - # flatten the probe payload to one plain string. - try: - _probe = json.loads(_probe_err) - _ts_scope_block = ( - f"{_probe.get('error', '')} Parameters schema: " - f"{json.dumps(_probe.get('parameters', {}), ensure_ascii=False)}. " - f"{_probe.get('hint', '')}" - ).strip() - except Exception: - _ts_scope_block = _probe_err - else: - function_name = _underlying - function_args = _underlying_args - else: - _ts_scope_block = ( - f"'{_underlying}' is not available in this session. " - "Use tool_search to find tools you can call." - ) - except Exception: - pass + function_name, function_args, _ts_scope_block = _unwrap_tool_search_call( + agent, function_name, function_args, flatten_probe=True + ) middleware_trace: list[dict[str, Any]] = [] _execution_blocked = False @@ -2089,14 +1939,17 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe tool_start_time = time.time() - if function_name == "todo_list": + if function_name != "delegate_task" and function_name in INLINE_TOOL_EXECUTORS: + # Agent-level tools that need live AIAgent state; table shared with invoke_tool. + inline_executor = INLINE_TOOL_EXECUTORS[function_name] + inline_ctx = InlineToolContext( + effective_task_id=effective_task_id, + tool_call_id=tool_call_id, + messages=messages, + ) + def _execute(next_args: dict) -> Any: - from tools.todo_tool import todo_tool as _todo_tool - return _todo_tool( - todos=next_args.get("todos"), - merge=next_args.get("merge", False), - store=agent._todo_store, - ) + return inline_executor(agent, next_args, inline_ctx) function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( agent, function_name=function_name, @@ -2109,290 +1962,7 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe )) tool_duration = time.time() - tool_start_time if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('todo_list', function_args, tool_duration, result=function_result)}") - elif function_name == "message_agent": - # Bot Mode teammate DM (tools/bot_mode_dm.py) — injected, not - # registered: only a canonical Bot Chat session carries the - # schema, and the tool re-gates on the session title itself. - def _execute(next_args: dict) -> Any: - from tools.bot_mode_dm import message_agent_tool as _message_agent_tool - return _message_agent_tool( - target=next_args.get("target", ""), - message=next_args.get("message", ""), - task_id=effective_task_id, - agent=agent, - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('message_agent', function_args, tool_duration, result=function_result)}") - elif function_name == "session_search": - def _execute(next_args: dict) -> Any: - session_db = agent._get_session_db_for_recall() - if not session_db: - from hermes_state import format_session_db_unavailable - return json.dumps({"success": False, "error": format_session_db_unavailable()}) - from tools.session_search_tool import session_search as _session_search - return _session_search( - query=next_args.get("query", ""), - role_filter=next_args.get("role_filter"), - limit=next_args.get("limit", 3), - session_id=next_args.get("session_id"), - around_message_id=next_args.get("around_message_id"), - window=next_args.get("window", 5), - sort=next_args.get("sort"), - detail=next_args.get("detail", "adaptive"), - db=session_db, - current_session_id=agent.session_id, - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('session_search', function_args, tool_duration, result=function_result)}") - elif function_name == "memory": - def _execute(next_args: dict) -> Any: - target = next_args.get("target", "memory") - operations = next_args.get("operations") - from tools.memory_tool import memory_tool as _memory_tool - result = _memory_tool( - action=next_args.get("action"), - target=target, - content=next_args.get("content"), - old_text=next_args.get("old_text"), - operations=operations, - store=agent._memory_store, - ) - # Mirror successful built-in memory writes to external - # providers. All gating/op-expansion lives behind the manager - # interface (MemoryManager.notify_memory_tool_write). - if agent._memory_manager: - agent._memory_manager.notify_memory_tool_write( - result, - next_args, - build_metadata=lambda: agent._build_memory_write_metadata( - task_id=effective_task_id, - tool_call_id=tool_call_id, - ), - ) - return result - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('memory', function_args, tool_duration, result=function_result)}") - elif function_name == "clarify": - def _execute(next_args: dict) -> Any: - from tools.clarify_tool import clarify_tool as _clarify_tool - return _clarify_tool( - question=next_args.get("question", ""), - choices=next_args.get("choices"), - multi_select=next_args.get("multi_select", False), - questions=next_args.get("questions"), - callback=agent.clarify_callback, - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('clarify', function_args, tool_duration, result=function_result)}") - elif function_name == "read_terminal": - def _execute(next_args: dict) -> Any: - from tools.read_terminal_tool import read_terminal_tool as _read_terminal_tool - return _read_terminal_tool( - start_line=next_args.get("start_line"), - count=next_args.get("count"), - callback=getattr(agent, "read_terminal_callback", None), - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('read_terminal', function_args, tool_duration, result=function_result)}") - elif function_name == "desktop_preview": - def _execute(next_args: dict) -> Any: - if (next_args.get("action") or "").strip() == "read": - from tools.read_preview_tool import read_preview_tool as _read_preview_tool - return _read_preview_tool( - start=next_args.get("start"), - count=next_args.get("count"), - callback=getattr(agent, "read_preview_callback", None), - ) - from tools.preview_tool import _handle_preview - return _handle_preview(next_args) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('desktop_preview', function_args, tool_duration, result=function_result)}") - elif function_name == "drive_preview": - def _execute(next_args: dict) -> Any: - from tools.drive_preview_tool import drive_preview_tool as _drive_preview_tool - return _drive_preview_tool( - action=next_args.get("action", ""), - ref=next_args.get("ref"), - selector=next_args.get("selector"), - text=next_args.get("text"), - key=next_args.get("key"), - submit=next_args.get("submit"), - amount=next_args.get("amount"), - to=next_args.get("to"), - limit=next_args.get("max"), - callback=getattr(agent, "drive_preview_callback", None), - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('drive_preview', function_args, tool_duration, result=function_result)}") - elif function_name == "annotate_preview": - def _execute(next_args: dict) -> Any: - from tools.annotate_preview_tool import annotate_preview_tool as _annotate_preview_tool - return _annotate_preview_tool( - action=next_args.get("action", "add"), - ref=next_args.get("ref"), - selector=next_args.get("selector"), - label=next_args.get("label"), - callback=getattr(agent, "drive_preview_callback", None), - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('annotate_preview', function_args, tool_duration, result=function_result)}") - elif function_name == "read_window_below": - def _execute(next_args: dict) -> Any: - from tools.read_window_tool import read_window_below_tool as _read_window_below_tool - return _read_window_below_tool( - callback=getattr(agent, "read_window_below_callback", None), - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('read_window_below', function_args, tool_duration, result=function_result)}") - elif function_name == "gui_tour": - def _execute(next_args: dict) -> Any: - from tools.tour_tool import tour_tool as _tour_tool - return _tour_tool( - action=next_args.get("action", ""), - surface=next_args.get("surface"), - selector=next_args.get("selector"), - title=next_args.get("title"), - text=next_args.get("text"), - side=next_args.get("side"), - steps=next_args.get("steps"), - step_index=next_args.get("step_index"), - callback=getattr(agent, "tour_callback", None), - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('gui_tour', function_args, tool_duration, result=function_result)}") - elif function_name == "setup_mcp": - def _execute(next_args: dict) -> Any: - from tools.setup_mcp_tool import setup_mcp_tool as _setup_mcp_tool - return _setup_mcp_tool( - server=next_args.get("server", ""), - action=next_args.get("action", "install"), - reason=next_args.get("reason", ""), - callback=getattr(agent, "setup_mcp_callback", None), - ) - function_result, function_args, middleware_trace, _execution_blocked, _execution_dispatched = _managed_values(_run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=getattr(tool_call, "id", "") or "", - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - )) - tool_duration = time.time() - tool_start_time - if agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {_get_cute_tool_message_impl('setup_mcp', function_args, tool_duration, result=function_result)}") + agent._vprint(f" {_get_cute_tool_message_impl(function_name, function_args, tool_duration, result=function_result)}") elif function_name == "delegate_task": _action_arg = str(function_args.get("action") or "").strip().lower() tasks_arg = function_args.get("tasks") @@ -2431,21 +2001,10 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe finally: agent._delegate_spinner = None tool_duration = time.time() - tool_start_time - cute_msg = _get_cute_tool_message_impl('delegate_task', function_args, tool_duration, result=_delegate_result) - if spinner: - spinner.stop(cute_msg) - elif agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {cute_msg}") + _finish_quiet_tool_spinner(agent, spinner, 'delegate_task', function_args, tool_duration, _delegate_result) elif agent._context_engine_tool_names and function_name in agent._context_engine_tool_names: # Context engine tools (lcm_grep, lcm_describe, lcm_expand, etc.) - spinner = None - if agent._should_emit_quiet_tool_messages(): - face = random.choice(KawaiiSpinner.get_waiting_faces()) - emoji = _get_tool_emoji(function_name) - display_args = _redact_tool_args_for_display(function_name, function_args) or function_args - preview = _build_tool_label(function_name, display_args) or function_name - spinner = KawaiiSpinner(f"{face} {emoji} {preview}", spinner_type='dots', print_fn=agent._print_fn) - spinner.start() + spinner = _start_quiet_tool_spinner(agent, function_name, function_args, gate=False) _ce_result = None try: def _execute(next_args: dict) -> Any: @@ -2466,22 +2025,11 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe logger.error("context_engine.handle_tool_call raised for %s: %s", function_name, tool_error, exc_info=True) finally: tool_duration = time.time() - tool_start_time - cute_msg = _get_cute_tool_message_impl(function_name, function_args, tool_duration, result=_ce_result) - if spinner: - spinner.stop(cute_msg) - elif agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {cute_msg}") + _finish_quiet_tool_spinner(agent, spinner, function_name, function_args, tool_duration, _ce_result) elif agent._memory_manager and agent._memory_manager.has_tool(function_name): # Memory provider tools (hindsight_retain, honcho_search, etc.) # These are not in the tool registry — route through MemoryManager. - spinner = None - if agent._should_emit_quiet_tool_messages() and agent._should_start_quiet_spinner(): - face = random.choice(KawaiiSpinner.get_waiting_faces()) - emoji = _get_tool_emoji(function_name) - display_args = _redact_tool_args_for_display(function_name, function_args) or function_args - preview = _build_tool_label(function_name, display_args) or function_name - spinner = KawaiiSpinner(f"{face} {emoji} {preview}", spinner_type='dots', print_fn=agent._print_fn) - spinner.start() + spinner = _start_quiet_tool_spinner(agent, function_name, function_args) _mem_result = None try: def _execute(next_args: dict) -> Any: @@ -2502,20 +2050,10 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe logger.error("memory_manager.handle_tool_call raised for %s: %s", function_name, tool_error, exc_info=True) finally: tool_duration = time.time() - tool_start_time - cute_msg = _get_cute_tool_message_impl(function_name, function_args, tool_duration, result=_mem_result) - if spinner: - spinner.stop(cute_msg) - elif agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {cute_msg}") - elif agent.quiet_mode: - spinner = None - if agent._should_emit_quiet_tool_messages() and agent._should_start_quiet_spinner(): - face = random.choice(KawaiiSpinner.get_waiting_faces()) - emoji = _get_tool_emoji(function_name) - display_args = _redact_tool_args_for_display(function_name, function_args) or function_args - preview = _build_tool_label(function_name, display_args) or function_name - spinner = KawaiiSpinner(f"{face} {emoji} {preview}", spinner_type='dots', print_fn=agent._print_fn) - spinner.start() + _finish_quiet_tool_spinner(agent, spinner, function_name, function_args, tool_duration, _mem_result) + else: + # Registry tools: post hook is owned by this executor (inner observer suppressed). + spinner = _start_quiet_tool_spinner(agent, function_name, function_args) if agent.quiet_mode else None _spinner_result = None try: def _execute(next_args: dict) -> Any: @@ -2579,9 +2117,8 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe agent.interrupt("keyboard interrupt") except Exception: pass - # Emit a tool result for THIS call and every remaining call in - # the batch before re-raising, so the assistant tool-call turn - # is never left without matching tool results (alternation). + # Emit results for THIS and every remaining call before re-raising so + # the tool-call turn keeps matching results (alternation). _append_cancelled_tool_results( messages, assistant_message.tool_calls[i - 1:], @@ -2593,84 +2130,8 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe logger.error("handle_function_call raised for %s: %s", function_name, tool_error, exc_info=True) finally: tool_duration = time.time() - tool_start_time - cute_msg = _get_cute_tool_message_impl(function_name, function_args, tool_duration, result=_spinner_result) - if spinner: - spinner.stop(cute_msg) - elif agent._should_emit_quiet_tool_messages(): - agent._vprint(f" {cute_msg}") - else: - try: - def _execute(next_args: dict) -> Any: - from model_tools import suppress_post_tool_call_hook - - with suppress_post_tool_call_hook(): - return _ra().handle_function_call( - function_name, - next_args, - effective_task_id, - tool_call_id=tool_call_id, - session_id=agent.session_id or "", - turn_id=getattr(agent, "_current_turn_id", "") or "", - api_request_id=getattr(agent, "_current_api_request_id", "") - or "", - enabled_tools=( - list(agent.valid_tool_names) - if agent.valid_tool_names - else None - ), - skip_pre_tool_call_hook=True, - skip_tool_request_middleware=True, - skip_tool_execution_middleware=True, - tool_request_middleware_trace=list(middleware_trace), - enabled_toolsets=getattr(agent, "enabled_toolsets", None), - disabled_toolsets=getattr(agent, "disabled_toolsets", None), - ) - - ( - function_result, - function_args, - middleware_trace, - _execution_blocked, - _execution_dispatched, - ) = _managed_values( - _run_agent_tool_execution_middleware( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - execute=_execute, - scope_block=_ts_scope_block, - display_index=i, - middleware_trace=middleware_trace, - ) - ) - except KeyboardInterrupt: - _emit_cancelled_terminal_post_tool_call( - agent, - function_name=function_name, - function_args=function_args, - effective_task_id=effective_task_id, - tool_call_id=tool_call_id, - start_time=tool_start_time, - middleware_trace=list(middleware_trace), - ) - try: - agent.interrupt("keyboard interrupt") - except Exception: - pass - # Emit a tool result for THIS call and every remaining call in - # the batch before re-raising (see interactive branch above). - _append_cancelled_tool_results( - messages, - assistant_message.tool_calls[i - 1:], - reason="keyboard interrupt", - ) - raise - except Exception as tool_error: - function_result = f"Error executing tool '{function_name}': {tool_error}" - logger.error("handle_function_call raised for %s: %s", function_name, tool_error, exc_info=True) - tool_duration = time.time() - tool_start_time + if agent.quiet_mode: + _finish_quiet_tool_spinner(agent, spinner, function_name, function_args, tool_duration, _spinner_result) _execution_timed_out = isinstance( function_result, (_ToolTimeoutResult, _ToolCancelledResult) @@ -2688,13 +2149,9 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe # Log tool errors to the persistent error log so [error] tags # in the UI always have a corresponding detailed entry on disk. _is_error_result, _ = _detect_tool_failure(function_name, function_result) - # The agent-runtime tools above (todo, session_search, memory, - # context-engine, memory-manager, clarify, delegate_task) are - # dispatched inline — they never reach handle_function_call, so the - # executor is the one that has to fire post_tool_call. For - # Every dispatch suppresses the inner handle_function_call observer so - # the executor owns one terminal event for this tool_call_id. This also - # prevents an abandoned timeout worker from reporting late success. + # Inline-dispatched runtime tools never reach handle_function_call, so the + # executor owns the one terminal post_tool_call per tool_call_id (the inner + # observer is suppressed); also stops an abandoned timeout worker reporting late. _executor_must_emit_post_hook = ( not _execution_blocked and not _execution_timed_out @@ -2726,10 +2183,8 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe else: logger.info("tool %s completed (%.2fs, %d chars)", function_name, tool_duration, _result_len) - # Track file-mutation outcome for the turn-end verifier. See - # the concurrent path for the rationale; both paths must feed - # the same state so the footer reflects every tool call in the - # turn, not just the parallel ones. + # Track file-mutation outcome for the turn-end verifier; both paths feed + # the same state so the footer reflects every tool call. if not _execution_blocked: try: agent._record_file_mutation_result( @@ -2748,84 +2203,35 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe logging.debug("Tool result (%d chars): %s", len(_log_result), _log_result) display_function_result = function_result - function_result = maybe_persist_tool_result( - content=function_result, - tool_name=function_name, - tool_use_id=tool_call_id, - env=get_active_env(effective_task_id), - config=_tool_budget, - ) if not _is_multimodal_tool_result(function_result) else function_result - _record_persisted_path_for_stub(agent, tool_call_id, function_result) - - # Discover subdirectory context files from tool arguments - subdir_hints = agent._subdirectory_hints.check_tool_call(function_name, function_args) - if subdir_hints: - if _is_multimodal_tool_result(function_result): - _append_subdir_hint_to_multimodal(function_result, subdir_hints) - else: - function_result += subdir_hints - - # Unwrap _multimodal dicts to an OpenAI-style content list - # (see parallel path for rationale). String results pass through. - _tool_content = agent._tool_result_content_for_active_model(function_name, function_result) - tool_message = make_tool_result_message( - function_name, - _tool_content, - tool_call_id, - effect_disposition="unknown" if _execution_timed_out else None, - ) - messages.append(tool_message) - risk_metadata = tool_message.get("_tool_output_risk") - if not _flush_session_db_after_tool_progress( + finalized = _append_finalized_tool_result( agent, messages, - stage=f"tool result {function_name}", - ): + function_name=function_name, + function_args=function_args, + function_result=function_result, + tool_call_id=tool_call_id, + effective_task_id=effective_task_id, + budget=_tool_budget, + effect_disposition="unknown" if _execution_timed_out else None, + ) + if finalized is None: return + function_result, _tool_message, risk_metadata = finalized - # UI completion/progress events are projections of the canonical tool - # row, never a competing in-memory authority. - if not _execution_blocked and agent.tool_progress_callback: - try: - agent.tool_progress_callback( - "tool.completed", function_name, None, None, - duration=tool_duration, is_error=_is_error_result, - result=display_function_result, - ) - except Exception as cb_err: - logging.debug("Tool progress callback error: %s", cb_err) - - if not _execution_blocked and agent.tool_complete_callback: - try: - display_args = ( - _redact_tool_args_for_display(function_name, function_args) - or function_args - ) - agent.tool_complete_callback( - tool_call_id, - function_name, - display_args, - display_function_result, - ) - except Exception as cb_err: - logging.debug("Tool complete callback error: %s", cb_err) - - if ( - risk_metadata is not None - and risk_metadata.get("risk") != "low" - and agent.tool_progress_callback - ): - try: - agent.tool_progress_callback( - "tool.output_risk", - function_name, - None, - None, - tool_call_id=tool_call_id, - risk_metadata=risk_metadata, - ) - except Exception as cb_err: - logging.debug("Tool output risk callback error: %s", cb_err) + if not _execution_blocked: + _emit_tool_completed_progress( + agent, function_name, + duration=tool_duration, is_error=_is_error_result, result=display_function_result, + ) + _emit_tool_complete_and_risk( + agent, + function_name=function_name, + function_args=function_args, + tool_call_id=tool_call_id, + result=display_function_result, + risk_metadata=risk_metadata, + blocked=_execution_blocked, + ) if not agent.quiet_mode and getattr(agent, "tool_progress_mode", "all") != "off": if agent.verbose_logging: @@ -2855,17 +2261,15 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe return break - # ── Per-turn aggregate budget enforcement ───────────────────────── - # Keep /steer pending until the final post-budget drain below. The model - # only receives this batch after all calls finish, and an early drain can - # be discarded when aggregate budget enforcement replaces a tool result. + # ── Per-turn aggregate budget enforcement ────────────────────────── + # Keep /steer pending until the post-budget drain: an early drain could be + # discarded when budget enforcement replaces a tool result. num_tools_seq = len(assistant_message.tool_calls) if finalize and num_tools_seq > 0: enforce_turn_budget(messages[-num_tools_seq:], env=get_active_env(effective_task_id), config=_tool_budget) - # ── /steer injection ────────────────────────────────────────────── - # See _execute_tool_calls_parallel for the rationale. Same hook, - # applied to sequential execution as well. + # ── /steer injection ──────────────────────────────────────────────── + # See the concurrent path for rationale. if finalize and num_tools_seq > 0: agent._apply_pending_steer_to_tool_results(messages, num_tools_seq) @@ -2873,27 +2277,13 @@ def execute_tool_calls_sequential(agent, assistant_message, messages: list, effe def execute_tool_calls_segmented(agent, assistant_message, messages: list, effective_task_id: str, api_call_count: int = 0, segments=None) -> None: - """Execute a mixed tool-call batch as ordered parallel/sequential segments. + """Execute a mixed batch as ordered parallel/sequential segments. - ``segments`` is the ``(kind, calls)`` plan from - ``_plan_tool_batch_segments``: maximal contiguous runs of parallel-safe - calls execute on the concurrent path, barrier calls on the sequential - path, strictly in the model's original call order. Because segments are - contiguous, every tool result is still appended one-per-call in emission - order and no call ever starts before an earlier barrier finishes — - identical ordering and side-effect boundaries to fully-sequential - execution, with I/O parallelism recovered inside the safe runs. - - Turn-end work (aggregate budget enforcement + /steer injection) is done - once here for the WHOLE batch; the per-segment executor calls run with - ``finalize=False`` so a multi-segment turn cannot multiply the budget or - truncate a steer marker. - - Interrupt semantics: each segment executor already checks - ``agent._interrupt_requested`` up front and appends a cancelled/skipped - result per call, so an interrupt during segment *k* drains segments - *k+1..n* without executing them while preserving one result per - tool_call_id. + ``segments`` is the ``(kind, calls)`` plan from ``_plan_tool_batch_segments``; + contiguous segments preserve per-call result order and barrier boundaries exactly + as fully-sequential execution. Turn-end work (budget + /steer) runs once here; + segment executors run with ``finalize=False``. Each segment executor checks the + interrupt flag up front, so an interrupt drains later segments with one result per call. """ from types import SimpleNamespace