diff --git a/tools/code_execution_env.py b/tools/code_execution_env.py index 00d40de390..7632eab67d 100644 --- a/tools/code_execution_env.py +++ b/tools/code_execution_env.py @@ -69,12 +69,9 @@ def _scrub_child_env(source_env, is_passthrough=None, is_windows=None): is_passthrough = is_env_passthrough if is_windows is None: is_windows = _IS_WINDOWS - scrubbed = {} - # Non-secret HERMES_* vars that no allowlist admits are dropped on purpose, - # but a script importing a repo module that reads one at import time would - # otherwise see it silently unset — log the drop once, pointing at the - # env_passthrough opt-in. + # Non-secret HERMES_* vars no allowlist admits are dropped on purpose; a script importing a + # repo module that reads one would see it silently unset — log the drop, point at the opt-in. _dropped_hermes = [] for k, v in source_env.items(): if is_passthrough(k): @@ -98,14 +95,11 @@ def _scrub_child_env(source_env, is_passthrough=None, is_windows=None): "env_passthrough in the skill/config so it passes by explicit opt-in.", len(_dropped_hermes), ", ".join(sorted(_dropped_hermes)), ) - - # delegate_task children are marked by a ContextVar, not os.environ, and the - # sandbox crosses a process boundary: bridge the marker and strip - # dispatcher-owned Kanban vars AFTER the scrub so an explicit passthrough - # cannot re-grant a delegated child the parent's board mutation capability. + # delegate_task children are marked by a ContextVar, not os.environ, and the sandbox crosses + # a process boundary: strip dispatcher-owned Kanban vars AFTER the scrub so an explicit + # passthrough cannot re-grant a delegated child the parent's board mutation capability. try: from agent.delegation_context import is_delegated_child_process_context, scrub_kanban_env - if is_delegated_child_process_context(): scrubbed = scrub_kanban_env(scrubbed) except Exception: @@ -121,10 +115,8 @@ def _build_child_env(*, rpc_endpoint: str, rpc_token: str, tmpdir: str, child_env["HERMES_RPC_SOCKET"] = rpc_endpoint child_env["HERMES_RPC_TOKEN"] = rpc_token child_env["PYTHONDONTWRITEBYTECODE"] = "1" - # Force UTF-8 stdio and default file encoding: on Windows sys.stdout is - # bound to the console code page (cp1252) and print("→") raises - # UnicodeEncodeError; PYTHONUTF8 also makes open()'s default UTF-8. - # Harmless belt-and-suspenders under a C/POSIX locale (minimal containers). + # Force UTF-8 stdio and default file encoding: on Windows sys.stdout is bound to the console + # code page (cp1252) and print("→") raises; harmless under a C/POSIX locale (containers). child_env["PYTHONIOENCODING"] = "utf-8" child_env["PYTHONUTF8"] = "1" # Only TZ reaches the child; HERMES_TIMEZONE is an internal setting. @@ -132,15 +124,11 @@ def _build_child_env(*, rpc_endpoint: str, rpc_token: str, tmpdir: str, if _tz_name: child_env["TZ"] = _tz_name child_env.pop("HERMES_TIMEZONE", None) - apply_subprocess_home_env(child_env) - # PYTHONPATH: the staging dir (hermes_tools.py lives there) must always be - # importable, even when project mode changes CWD. Hermes's own root is - # added ONLY when the child runs in Hermes's Python environment — exposing - # Hermes's site-packages to an external project interpreter can mix - # incompatible compiled extensions (3.12 NumPy under a 3.9 venv). Inherited - # Hermes-owned entries (PYTHONPATH passes the scrub) are stripped first so - # they never shadow the child's sys.path. + # PYTHONPATH: the staging dir (hermes_tools.py) must always be importable even when project + # mode changes CWD. Hermes's root is added ONLY when the child runs in Hermes's Python env — + # exposing Hermes's site-packages to an external interpreter can mix incompatible compiled + # extensions (3.12 NumPy under a 3.9 venv). Inherited Hermes-owned entries are stripped first. from tools.environments.local import _strip_hermes_owned_pythonpath _strip_hermes_owned_pythonpath(child_env) _existing_pp = child_env.get("PYTHONPATH", "") @@ -148,8 +136,7 @@ def _build_child_env(*, rpc_endpoint: str, rpc_token: str, tmpdir: str, if _uses_hermes_python_environment(child_python): _pp_parts.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) elif child_python not in _external_env_logged: - # Import behavior changes silently otherwise — surface it once per - # interpreter path so "import hermes_constants fails" is diagnosable. + # Surface once per interpreter so "import hermes_constants fails" is diagnosable. _external_env_logged.add(child_python) logger.info("execute_code: child interpreter %s is outside the Hermes " "environment; hermes root omitted from PYTHONPATH", child_python) @@ -159,9 +146,8 @@ def _build_child_env(*, rpc_endpoint: str, rpc_token: str, tmpdir: str, return child_env -# Interpreter-probe caches: success-only dicts (FIFO-evicted at the cap) rather -# than lru_cache — a transient probe failure (fork pressure, 5s timeout on a -# loaded host) must not stick for the process lifetime. +# Interpreter-probe caches: success-only dicts (FIFO-evicted at the cap) rather than lru_cache — +# a transient probe failure (fork pressure, 5s timeout) must not stick for the process lifetime. _PROBE_CACHE_MAX = 32 _usable_python_cache: dict = {} _python_prefix_cache: dict = {} @@ -181,7 +167,6 @@ def _probe_python(python_path: str, code: str, *, text: bool = False): """Run ``python_path -c code``; None if missing, unspawnable, or past the 5s timeout.""" try: from agent.delegation_context import delegated_child_subprocess_env - return subprocess.run( [python_path, "-c", code], timeout=5, capture_output=True, text=text, creationflags=subprocess.CREATE_NO_WINDOW if _IS_WINDOWS else 0, @@ -218,13 +203,9 @@ def _python_environment_prefix(python_path: str) -> str: def _uses_hermes_python_environment(python_path: str) -> bool: - """Whether *python_path* belongs to Hermes's active Python environment. - - Short-circuits when it IS the running interpreter (by path or realpath) so - no probe runs on the default strict path and a flaky probe of - sys.executable can never drop the hermes root; the realpath leg also covers - venvs whose bin/python resolves to the same binary (``uv run``). - """ + """Whether *python_path* belongs to Hermes's active Python environment. Short-circuits when + it IS the running interpreter (by path or realpath — covers ``uv run`` venvs) so no probe + runs on the default strict path and a flaky probe can never drop the hermes root.""" if python_path == sys.executable or os.path.realpath(python_path) == os.path.realpath(sys.executable): return True return _python_environment_prefix(python_path) == os.path.realpath(sys.prefix) @@ -236,7 +217,6 @@ def _resolve_child_python(mode: str) -> str: 3.8+ probe, else ``sys.executable``.""" if mode != "project": return sys.executable - subdir, exe_names = ("Scripts", ("python.exe", "python3.exe")) if _IS_WINDOWS else ("bin", ("python", "python3")) for var in ("VIRTUAL_ENV", "CONDA_PREFIX"): root = os.environ.get(var, "").strip() @@ -251,24 +231,18 @@ def _resolve_child_python(mode: str) -> str: logger.info("execute_code: skipping %s=%s (Python version < 3.8 or broken). " "Using sys.executable instead.", var, candidate) return sys.executable - return sys.executable def _resolve_child_cwd(mode: str, staging_dir: str, task_id: str = "") -> str: - """Working directory for the child. - - Strict mode: the staging dir. Project mode mirrors the terminal/file-tool - ladder so every file-writing path in a session agrees: the session's cwd - record (its `cd` state) → registered ``session.cwd.set`` override → - TERMINAL_CWD → os.getcwd() → staging dir (never Popen on a missing cwd). - """ + """Child cwd. Strict: the staging dir. Project mirrors the terminal/file-tool ladder so every + file-writing path agrees: session cwd record (`cd` state) → registered ``session.cwd.set`` + override → TERMINAL_CWD → os.getcwd() → staging dir (never Popen on a missing cwd).""" if mode != "project": return staging_dir if task_id: try: from tools.terminal_tool import get_session_cwd - recorded = get_session_cwd(task_id) except Exception: recorded = None @@ -276,14 +250,12 @@ def _resolve_child_cwd(mode: str, staging_dir: str, task_id: str = "") -> str: return recorded try: from tools.file_tools import _registered_task_cwd_override - session_cwd = _registered_task_cwd_override(task_id) except Exception: session_cwd = None if session_cwd and os.path.isdir(session_cwd): return session_cwd from agent.runtime_cwd import scope_terminal_cwd - raw = scope_terminal_cwd().strip() if raw: expanded = os.path.expanduser(raw) diff --git a/tools/code_execution_rpc.py b/tools/code_execution_rpc.py index 9e242e422f..a3191a27d8 100644 --- a/tools/code_execution_rpc.py +++ b/tools/code_execution_rpc.py @@ -27,15 +27,12 @@ _TERMINAL_BLOCKED_PARAMS = {"background", "pty", "notify", "notify_on_complete", def _default_dispatch(task_id): from model_tools import handle_function_call - return lambda tool_name, tool_args: handle_function_call(tool_name, tool_args, task_id=task_id) def _rpc_token_ok(request: dict, rpc_token: str) -> bool: - """Constant-time token check; an empty server token fails closed. - - Compared as bytes: compare_digest raises TypeError on a non-ASCII str, and - the token comes from script-supplied JSON.""" + """Constant-time token check; an empty server token fails closed. Compared as bytes: + compare_digest raises TypeError on a non-ASCII str, and the token is script-supplied JSON.""" return bool(rpc_token) and secrets.compare_digest( str(request.get("token") or "").encode(), rpc_token.encode() ) @@ -44,24 +41,19 @@ def _rpc_token_ok(request: dict, rpc_token: str) -> bool: def _handle_rpc_request(request: dict, *, allowed_tools: frozenset, tool_call_counter: list, max_tool_calls: int, dispatch, tool_call_log: list, call_start: float, where: str) -> str: - """Enforce allow-list + budget, then dispatch one authenticated request. - - Only a dispatched call consumes budget and is logged; refusals are free. - """ + """Enforce allow-list + budget, then dispatch one authenticated request. Only a dispatched + call consumes budget and is logged; refusals are free.""" tool_name = request.get("tool", "") tool_args = request.get("args", {}) - if tool_name not in allowed_tools: return tool_error(f"Tool '{tool_name}' is not available in execute_code. " f"Available: {', '.join(sorted(allowed_tools))}") if tool_call_counter[0] >= max_tool_calls: return tool_error(f"Tool call limit reached ({max_tool_calls}). " "No more tool calls allowed in this execution.") - if tool_name == "terminal" and isinstance(tool_args, dict): for param in _TERMINAL_BLOCKED_PARAMS: tool_args.pop(param, None) - # Silence handler status prints so they don't leak into the CLI spinner. try: with thread_scoped_silence(): @@ -69,7 +61,6 @@ def _handle_rpc_request(request: dict, *, allowed_tools: frozenset, tool_call_co except Exception as exc: logger.error("Tool call failed in %s: %s", where, exc, exc_info=True) result = tool_error(str(exc)) - tool_call_counter[0] += 1 tool_call_log.append({"tool": tool_name, "args_preview": str(tool_args)[:80], "duration": round(time.monotonic() - call_start, 2)}) @@ -79,18 +70,13 @@ def _handle_rpc_request(request: dict, *, allowed_tools: frozenset, tool_call_co def _rpc_server_loop(server_sock: socket.socket, task_id: str, tool_call_log: list, tool_call_counter: list, max_tool_calls: int, allowed_tools: frozenset, stop_event: threading.Event, rpc_token: str, dispatch=None): - """Accept one client and serve newline-delimited JSON requests until it - disconnects, idles 300s, or the call limit is reached. - - ``tool_call_counter`` is a mutable ``[int]`` so the thread can increment it. - ``dispatch`` overrides how an allowed, budgeted call runs: per-call - sandboxes use the default (the thread already carries the cell's context), - while session kernels pass a dispatcher that rebinds each call to the - CURRENT cell's authority — the serving thread outlives many cells there. + """Accept one client and serve newline-delimited JSON requests until it disconnects, idles + 300s, or the call limit is reached. ``tool_call_counter`` is a mutable ``[int]``. ``dispatch`` + overrides how an allowed, budgeted call runs: per-call sandboxes use the default (the thread + carries the cell's context); session kernels rebind each call to the CURRENT cell's authority. """ if dispatch is None: dispatch = _default_dispatch(task_id) - conn = None try: server_sock.settimeout(0.05) @@ -103,7 +89,6 @@ def _rpc_server_loop(server_sock: socket.socket, task_id: str, tool_call_log: li if conn is None: return conn.settimeout(300) - buf = b"" while True: try: @@ -113,13 +98,11 @@ def _rpc_server_loop(server_sock: socket.socket, task_id: str, tool_call_log: li if not chunk: break buf += chunk - while b"\n" in buf: line, buf = buf.split(b"\n", 1) line = line.strip() if not line: continue - call_start = time.monotonic() try: request = json.loads(line.decode()) @@ -135,7 +118,6 @@ def _rpc_server_loop(server_sock: socket.socket, task_id: str, tool_call_log: li call_start=call_start, where="sandbox", ) conn.sendall((resp + "\n").encode()) - except socket.timeout: logger.debug("RPC listener socket timeout") except OSError as e: @@ -151,15 +133,11 @@ def _rpc_server_loop(server_sock: socket.socket, task_id: str, tool_call_log: li def _rpc_poll_loop(env, rpc_dir: str, task_id: str, tool_call_log: list, tool_call_counter: list, max_tool_calls: int, allowed_tools: frozenset, stop_event: threading.Event, rpc_token: str): - """Poll the remote filesystem for request files and answer them. - - Runs in a background thread; each ``env.execute()`` is an independent - process, so this is safe alongside the script-execution thread. Malformed - or unauthorized requests are removed without a response. - """ + """Poll the remote filesystem for request files and answer them. Background thread; each + ``env.execute()`` is an independent process, so this is safe alongside the script-execution + thread. Malformed or unauthorized requests are removed without a response.""" dispatch = _default_dispatch(task_id) poll_interval = 0.1 - quoted_rpc_dir = shlex.quote(rpc_dir) while not stop_event.is_set(): try: @@ -168,16 +146,13 @@ def _rpc_poll_loop(env, rpc_dir: str, task_id: str, tool_call_log: list, tool_ca if not output: stop_event.wait(poll_interval) continue - req_files = sorted( f for f in (line.strip() for line in output.split("\n")) if f and not f.endswith(".tmp") and "/req_" in f ) - for req_file in req_files: if stop_event.is_set(): break - call_start = time.monotonic() quoted_req_file = shlex.quote(req_file) read_result = env.execute(f"cat {quoted_req_file}", cwd="/", timeout=10) @@ -191,13 +166,11 @@ def _rpc_poll_loop(env, rpc_dir: str, task_id: str, tool_call_log: list, tool_ca logger.debug("Unauthorized RPC request in %s", req_file) env.execute(f"rm -f {quoted_req_file}", cwd="/", timeout=5) continue - tool_result = _handle_rpc_request( request, allowed_tools=allowed_tools, tool_call_counter=tool_call_counter, max_tool_calls=max_tool_calls, dispatch=dispatch, tool_call_log=tool_call_log, call_start=call_start, where="remote sandbox", ) - # Write the response atomically (tmp + rename) via echo piping — # Modal doesn't reliably deliver stdin_data to chained commands. quoted_res_file = shlex.quote(f"{rpc_dir}/res_{request.get('seq', 0):06d}") @@ -208,10 +181,8 @@ def _rpc_poll_loop(env, rpc_dir: str, task_id: str, tool_call_log: list, tool_ca cwd="/", timeout=60, ) env.execute(f"rm -f {quoted_req_file}", cwd="/", timeout=5) - except Exception as e: if not stop_event.is_set(): logger.debug("RPC poll error: %s", e, exc_info=True) - if not stop_event.is_set(): stop_event.wait(poll_interval) diff --git a/tools/code_execution_tool.py b/tools/code_execution_tool.py index dab099d081..99c97aca86 100644 --- a/tools/code_execution_tool.py +++ b/tools/code_execution_tool.py @@ -1,19 +1,14 @@ #!/usr/bin/env python3 -""" -Code Execution Tool -- Programmatic Tool Calling (PTC) +"""Code Execution Tool -- Programmatic Tool Calling (PTC). -Lets the LLM write a Python script that calls Hermes tools via RPC, collapsing -multi-step tool chains into a single inference turn. Only the script's stdout -is returned to the LLM; intermediate tool results never enter the context. - -Two transports: the local backend runs a persistent per-conversation session -kernel (tools/code_kernel.py) talking to the parent's RPC thread over a Unix -domain socket (loopback TCP on Windows); remote backends (Docker/SSH/Modal/...) -run a remote session kernel (tools/code_kernel_remote.py), falling open to a -per-call script ship, with tool calls as request files that a polling thread on -the parent reads via env.execute(). Sibling modules: tools/code_execution_env.py -(env scrubbing, interpreter/cwd resolution) and tools/code_execution_rpc.py (RPC -servers). Remote execution requires Python 3 in the terminal backend. +The LLM writes a Python script that calls Hermes tools via RPC, collapsing +multi-step tool chains into one inference turn; only the script's stdout returns +to the LLM. Local backend: a persistent per-conversation session kernel +(tools/code_kernel.py) over a Unix socket (loopback TCP on Windows). Remote +backends: a remote session kernel (tools/code_kernel_remote.py) falling open to a +per-call script ship, tool calls as request files polled via env.execute() +(needs Python 3 on the backend). Siblings: tools/code_execution_env.py (env +scrubbing, interpreter/cwd), tools/code_execution_rpc.py (RPC servers). """ import base64 @@ -36,38 +31,19 @@ from tools.registry import registry, tool_error # Env/interpreter resolution and RPC servers live in sibling modules; re-exported # here so `from tools.code_execution_tool import X` / patch() targets keep working. from tools.code_execution_env import ( # noqa: F401 - _SAFE_ENV_PREFIXES, - _SECRET_SUBSTRINGS, - _HERMES_CHILD_ALLOWED, - _WINDOWS_ESSENTIAL_ENV_VARS, - _scrub_child_env, - _build_child_env, - _PROBE_CACHE_MAX, - _usable_python_cache, - _python_prefix_cache, - _external_env_logged, - _cache_probe_result, - _is_usable_python, - _probe_python, - _python_environment_prefix, - _uses_hermes_python_environment, - _resolve_child_python, - _resolve_child_cwd, -) -from tools.code_execution_rpc import ( # noqa: F401 - _TERMINAL_BLOCKED_PARAMS, - _rpc_server_loop, - _rpc_poll_loop, + _SAFE_ENV_PREFIXES, _SECRET_SUBSTRINGS, _HERMES_CHILD_ALLOWED, _WINDOWS_ESSENTIAL_ENV_VARS, + _scrub_child_env, _build_child_env, _PROBE_CACHE_MAX, _usable_python_cache, _python_prefix_cache, + _external_env_logged, _cache_probe_result, _is_usable_python, _probe_python, + _python_environment_prefix, _uses_hermes_python_environment, _resolve_child_python, _resolve_child_cwd, ) +from tools.code_execution_rpc import _TERMINAL_BLOCKED_PARAMS, _rpc_server_loop, _rpc_poll_loop # noqa: F401 logger = logging.getLogger(__name__) -# Loopback TCP replaces AF_UNIX on Windows, so execute_code is available on -# every platform Hermes itself runs on. +# Loopback TCP replaces AF_UNIX on Windows, so execute_code runs on every platform Hermes does. SANDBOX_AVAILABLE = True -# Tools allowed inside the sandbox; the intersection with the session's -# enabled tools determines which stubs are generated. +# Tools allowed inside the sandbox; ∩ the session's enabled tools decides which stubs are generated. SANDBOX_ALLOWED_TOOLS = frozenset([ "web_search", "web_extract", "read_file", "write_file", "search_files", "patch", "terminal", ]) @@ -85,11 +61,9 @@ MAX_SPILLED_STDOUT_BYTES = 5_000_000 def _truncate_stdout_text(stdout_text: str) -> Tuple[str, Dict[str, Any]]: """Cap stdout by bytes (40% head / 60% tail) with explicit truncation metadata. - The agent receives execute_code results as JSON; a textual marker can be - missed or re-truncated by a client layer, so byte counts ride alongside it. - The omitted middle is not discarded: the complete output is spilled to - cache/exec and the result carries the path — the same recover-don't-rerun - pattern as web_extract's cache/web full-text store. + Byte counts ride alongside the textual marker because a client layer can miss + or re-truncate the marker. The omitted middle is spilled to cache/exec and the + result carries the path (recover-don't-rerun, as web_extract's cache/web). """ stdout_bytes = stdout_text.encode("utf-8", errors="replace") total = len(stdout_bytes) @@ -98,7 +72,6 @@ def _truncate_stdout_text(stdout_text: str) -> Tuple[str, Dict[str, Any]]: "stdout_truncated": False, "stdout_bytes_captured": total, "stdout_bytes_total": total, "stdout_bytes_omitted": 0, } - head_bytes = int(MAX_STDOUT_BYTES * 0.4) omitted = total - MAX_STDOUT_BYTES text = ( @@ -129,18 +102,13 @@ def _truncate_stdout_text(stdout_text: str) -> Tuple[str, Dict[str, Any]]: def _spill_full_stdout(stdout_text: str) -> Optional[str]: - """Write full stdout to cache/exec; return its path (None on failure). - - Best-effort by design — truncated inline output is still returned when - storage fails. Files are keyed by content digest so identical reruns - coalesce; the directory rides the same remote bind-mount list as - cache/web (credential_files._CACHE_DIRS) if present there. - """ + """Write full stdout to cache/exec; return its path (None on failure — best-effort, + the truncated inline output is still returned). Keyed by content digest so identical + reruns coalesce; the dir rides the cache/web remote bind-mount list (credential_files).""" try: import hashlib from hermes_constants import get_hermes_dir from tools.spill_safety import write_text_exclusive - if len(stdout_text) > MAX_SPILLED_STDOUT_BYTES: stdout_text = (stdout_text[:MAX_SPILLED_STDOUT_BYTES] + f"\n\n[... spill capped at {MAX_SPILLED_STDOUT_BYTES:,} bytes ...]") @@ -161,7 +129,6 @@ def check_sandbox_requirements() -> bool: return False try: from tools.terminal_tool import _check_vercel_sandbox_requirements, _get_env_config - config = _get_env_config() except Exception: logger.debug("Could not resolve terminal config for execute_code availability", exc_info=True) @@ -175,8 +142,7 @@ def check_sandbox_requirements() -> bool: # hermes_tools.py code generator # --------------------------------------------------------------------------- -# Per-tool stub templates: (signature, docstring, args_dict_expr); the -# args_dict_expr builds the JSON payload sent over the RPC channel. +# Per-tool stub templates: (signature, docstring, args_dict_expr — the JSON payload sent over RPC). _TOOL_STUBS = { "web_search": ( "query: str, limit: int = 5", @@ -227,10 +193,9 @@ def _missing_hermes_tools_import_hint(m, enabled_tools) -> str: "else, use the normal tool call instead of execute_code.") -# (regex, formatter(match, enabled_tools)) — first match wins. Production -# mining (state.db) ranked these as the top execute_code failure classes: -# hermes_tools import misuse, importing the built-in helpers, treating tool -# results as strings, importing third-party packages absent from the sandbox. +# (regex, formatter(match, enabled_tools)) — first match wins. Production mining (state.db) ranked +# these as the top execute_code failure classes: hermes_tools import misuse, importing the built-in +# helpers, treating tool results as strings, importing third-party packages absent from the sandbox. _FAILURE_HINT_RULES = ( (r"cannot import name '(\w+)' from 'hermes_tools'", _missing_hermes_tools_import_hint), (r"NameError: name '(json_parse|shell_quote|retry)' is not defined", @@ -242,8 +207,7 @@ _FAILURE_HINT_RULES = ( "terminal() with the project venv's python instead."), (r"TypeError: string indices must be integers|AttributeError: 'str' object has no attribute 'get'", lambda m, _: "Tool functions in the sandbox return DICTS (already parsed) — " - "do not json.loads() them or index them like strings. " - "Example: read_file(path)['content']."), + "do not json.loads() them or index them like strings. Example: read_file(path)['content']."), ) @@ -265,19 +229,13 @@ def _sandbox_failure_hint(stderr_text: str, enabled_tools=None) -> Optional[str] def generate_hermes_tools_module(enabled_tools: List[str], transport: str = "uds") -> str: - """Source of the hermes_tools.py stub module for tools in both - SANDBOX_ALLOWED_TOOLS and *enabled_tools*. ``transport``: ``"uds"`` (local - socket client) or ``"file"`` (file-based RPC client for remote backends).""" + """Source of the hermes_tools.py stub module for SANDBOX_ALLOWED_TOOLS ∩ *enabled_tools*. + ``transport``: ``"uds"`` (local socket client) or ``"file"`` (file RPC, remote backends).""" header = _FILE_TRANSPORT_HEADER if transport == "file" else _UDS_TRANSPORT_HEADER - stubs = [] - for name in sorted(SANDBOX_ALLOWED_TOOLS & set(enabled_tools)): - sig, doc, args_expr = _TOOL_STUBS[name] - stubs.append( - f"def {name}({sig}):\n" - f" {doc}\n" - f" return _call({name!r}, {args_expr})\n" - ) - return header + "\n".join(stubs) + return header + "\n".join( + f"def {name}({sig}):\n {doc}\n return _call({name!r}, {args_expr})\n" + for name, (sig, doc, args_expr) in sorted(_TOOL_STUBS.items()) if name in set(enabled_tools) + ) # ---- Shared helpers section (embedded in both transport headers) ---------- @@ -480,41 +438,33 @@ def _call(tool_name, args): # --------------------------------------------------------------------------- def _get_or_create_env(task_id: str): - """Return ``(env, env_type)`` — the same environment the terminal and file - tools use for *task_id*, created on first use (same double-checked - per-task lock pattern as file_tools._get_file_ops).""" + """``(env, env_type)`` — the environment the terminal/file tools share for *task_id*, created on + first use (same double-checked per-task lock pattern as file_tools._get_file_ops).""" from tools.terminal_tool import ( _active_environments, _env_lock, _create_environment, _get_env_config, _last_activity, _start_cleanup_thread, _creation_locks, _creation_locks_lock, _task_env_overrides, _resolve_container_task_id, _resolve_task_host_cwd, _is_container_backend, _select_image, _ssh_config_from_config, ) - effective_task_id = _resolve_container_task_id(task_id) - def _cached(): with _env_lock: env = _active_environments.get(effective_task_id) if env is not None: _last_activity[effective_task_id] = time.time() return env - env = _cached() if env is not None: return env, _get_env_config()["env_type"] - with _creation_locks_lock: task_lock = _creation_locks.setdefault(effective_task_id, threading.Lock()) - with task_lock: env = _cached() if env is not None: return env, _get_env_config()["env_type"] - config = _get_env_config() env_type = config["env_type"] overrides = _task_env_overrides.get(effective_task_id, {}) - container_config = None if _is_container_backend(env_type): container_config = { @@ -527,25 +477,19 @@ def _get_or_create_env(task_id: str): "docker_run_as_host_user": config.get("docker_run_as_host_user", False), "docker_network": config.get("docker_network", True), } - logger.info("Creating new %s environment for execute_code task %s...", env_type, effective_task_id[:8]) env = _create_environment( - env_type=env_type, - image=_select_image(env_type, overrides, config), - cwd=overrides.get("cwd") or config["cwd"], - timeout=config["timeout"], + env_type=env_type, image=_select_image(env_type, overrides, config), + cwd=overrides.get("cwd") or config["cwd"], timeout=config["timeout"], ssh_config=_ssh_config_from_config(config) if env_type == "ssh" else None, container_config=container_config, local_config={"persistent": config.get("local_persistent", False)} if env_type == "local" else None, - task_id=effective_task_id, - host_cwd=_resolve_task_host_cwd(config, task_id), + task_id=effective_task_id, host_cwd=_resolve_task_host_cwd(config, task_id), ) - with _env_lock: _active_environments[effective_task_id] = env _last_activity[effective_task_id] = time.time() - _start_cleanup_thread() logger.info("%s environment ready for execute_code task %s", env_type, effective_task_id[:8]) @@ -553,9 +497,8 @@ def _get_or_create_env(task_id: str): def _ship_file_to_remote(env, remote_path: str, content: str) -> None: - """Write *content* to *remote_path* via ``echo … | base64 -d`` — some - backends (Modal) don't reliably deliver stdin_data to chained commands, and - base64 is shell-safe inside single quotes.""" + """Write *content* to *remote_path* via ``echo … | base64 -d`` — some backends (Modal) don't + reliably deliver stdin_data to chained commands; base64 is shell-safe inside single quotes.""" encoded = base64.b64encode(content.encode("utf-8")).decode("ascii") env.execute(f"echo '{encoded}' | base64 -d > {shlex.quote(remote_path)}", cwd="/", timeout=30) @@ -578,38 +521,29 @@ def _env_temp_dir(env: Any) -> str: def _format_interrupted_output(stdout_text: str) -> str: """Append an interruption marker without guessing who caused it.""" from tools.interrupt import get_interrupt_reason - reason = get_interrupt_reason() marker = f"[execution interrupted — {reason}]" if reason else "[execution interrupted]" return f"{stdout_text}\n{marker}" if stdout_text else marker def _clean_output(stdout_text: str) -> Tuple[str, Dict[str, Any]]: - """Shared output pipeline: byte-cap (with spill), ANSI strip, secret redaction. - - code_file=True: execution output often echoes source/config — skip the - ENV/JSON/f-string-template false positives while still masking real credentials. - """ + """Shared output pipeline: byte-cap (with spill), ANSI strip, secret redaction. code_file=True: + output often echoes source/config — skip ENV/JSON/f-string false positives, still mask credentials.""" from tools.ansi_strip import strip_ansi from agent.redact import redact_sensitive_text - stdout_text, metadata = _truncate_stdout_text(stdout_text) return redact_sensitive_text(strip_ansi(stdout_text), code_file=True), metadata def _with_timeout_notice(stdout_text: str, timeout_msg: str) -> str: - """Put the timeout message in the output too — an empty result makes models - answer as if nothing happened, and the gateway drops empty replies.""" + """Timeout message goes in the output too — an empty result makes models answer as if + nothing happened, and the gateway drops empty replies.""" return stdout_text + f"\n\n⏰ {timeout_msg}" if stdout_text else f"⏰ {timeout_msg}" def _error_result(error: str, *, tool_calls_made: int = 0, duration: float = 0) -> str: - return json.dumps({ - "status": "error", - "error": error, - "tool_calls_made": tool_calls_made, - "duration_seconds": duration, - }, ensure_ascii=False) + return json.dumps({"status": "error", "error": error, "tool_calls_made": tool_calls_made, + "duration_seconds": duration}, ensure_ascii=False) def _remote_failure(exc: BaseException, exec_start: float, tool_calls_made: int) -> str: @@ -622,11 +556,28 @@ def _remote_failure(exc: BaseException, exec_start: float, tool_calls_made: int) _REMOTE_EXIT_STATUS = {124: "timeout", 130: "interrupted"} +def _remote_result(status: str, raw_stdout: str, exec_start: float, fields: Dict[str, Any], + kernel: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: + """Common remote reply shape: status, cleaned output, *fields*, duration, optional kernel + info, then truncation metadata (key order is part of the result contract).""" + stdout_text, stdout_metadata = _clean_output(raw_stdout) + result: Dict[str, Any] = {"status": status, "output": stdout_text, **fields, + "duration_seconds": round(time.monotonic() - exec_start, 2)} + if kernel is not None: + result["kernel"] = kernel + result.update(stdout_metadata) + return result + + +def _apply_timeout(result: Dict[str, Any], timeout_msg: str) -> None: + result["error"] = timeout_msg + result["output"] = _with_timeout_notice(result["output"], timeout_msg) + + def _finish_remote_kernel_result(kernel_result: Dict[str, Any], *, timeout: int, exec_start: float) -> str: - """Post-process a remote-kernel cell result into the tool's JSON reply. - Timeout messaging mirrors the local kernel contract (kernel killed, state - lost, next call fresh).""" + """Post-process a remote-kernel cell result into the tool's JSON reply. Timeout messaging + mirrors the local kernel contract (kernel killed, state lost, next call fresh).""" stdout_text = kernel_result.get("stdout", "") or "" stderr_text = kernel_result.get("stderr", "") or "" traceback_text = kernel_result.get("traceback", "") or "" @@ -634,42 +585,30 @@ def _finish_remote_kernel_result(kernel_result: Dict[str, Any], *, # Same joining shape as the local kernel: stderr and traceback ride in # the output under one marker so the model sees the failure inline. stdout_text = stdout_text + "\n--- stderr ---\n" + stderr_text + traceback_text - - stdout_text, stdout_metadata = _clean_output(stdout_text) - result: Dict[str, Any] = { - "status": kernel_result.get("status", "error"), - "output": stdout_text, - "tool_calls_made": kernel_result.get("tool_calls_made", 0), - "duration_seconds": round(time.monotonic() - exec_start, 2), - "kernel": kernel_result.get("kernel", {"remote": True}), - } - result.update(stdout_metadata) - + result = _remote_result(kernel_result.get("status", "error"), stdout_text, exec_start, + {"tool_calls_made": kernel_result.get("tool_calls_made", 0)}, + kernel=kernel_result.get("kernel", {"remote": True})) if result["status"] == "timeout": - timeout_msg = (f"Cell timed out after {timeout}s; the remote session kernel was " - "killed and its state was lost. The next call starts fresh.") - result["error"] = timeout_msg - result["output"] = _with_timeout_notice(stdout_text, timeout_msg) + _apply_timeout(result, f"Cell timed out after {timeout}s; the remote session kernel was " + "killed and its state was lost. The next call starts fresh.") elif result["status"] == "error" and kernel_result.get("error"): result["error"] = kernel_result["error"] return json.dumps(result, ensure_ascii=False) def _sandbox_tools_for(enabled_tools: Optional[List[str]]) -> frozenset: - """Tools the sandbox may call: the session's enabled set ∩ SANDBOX_ALLOWED_TOOLS, - or every sandbox tool when the intersection is empty.""" + """Enabled ∩ SANDBOX_ALLOWED_TOOLS, or every sandbox tool when the intersection is empty.""" return frozenset(SANDBOX_ALLOWED_TOOLS & set(enabled_tools or ())) or SANDBOX_ALLOWED_TOOLS def _run_remote_per_call(env, env_type: str, code: str, effective_task_id: str, sandbox_tools: frozenset, *, timeout: int, max_tool_calls: int, exec_start: float) -> str: - """Per-call script ship: stage hermes_tools.py + script.py in a fresh remote - sandbox dir, serve file-RPC from a polling thread, run, clean up.""" + """Per-call script ship: stage hermes_tools.py + script.py in a fresh remote sandbox dir, + serve file-RPC from a polling thread, run, clean up.""" sandbox_dir = f"{_env_temp_dir(env)}/hermes_exec_{uuid.uuid4().hex[:12]}" quoted_sandbox_dir = shlex.quote(sandbox_dir) quoted_rpc_dir = shlex.quote(f"{sandbox_dir}/rpc") - tool_call_log: list = [] tool_call_counter = [0] stop_event = threading.Event() @@ -680,7 +619,6 @@ def _run_remote_per_call(env, env_type: str, code: str, effective_task_id: str, _ship_file_to_remote(env, f"{sandbox_dir}/hermes_tools.py", generate_hermes_tools_module(list(sandbox_tools), transport="file")) _ship_file_to_remote(env, f"{sandbox_dir}/script.py", code) - # Wrapped so the thread inherits the turn's approval context + callbacks # (tools.thread_context) — else sandbox RPC tool calls lose approval routing. rpc_thread = threading.Thread( @@ -690,14 +628,12 @@ def _run_remote_per_call(env, env_type: str, code: str, effective_task_id: str, daemon=True, ) rpc_thread.start() - env_prefix = (f"HERMES_RPC_DIR={quoted_rpc_dir} " f"HERMES_RPC_TOKEN={shlex.quote(rpc_token)} " f"PYTHONDONTWRITEBYTECODE=1") tz = os.getenv("HERMES_TIMEZONE", "").strip() if tz: env_prefix += f" TZ={shlex.quote(tz)}" - logger.info("Executing code on %s backend (task %s)...", env_type, effective_task_id[:8]) script_result = env.execute(f"cd {quoted_sandbox_dir} && {env_prefix} python3 script.py", timeout=timeout) @@ -715,26 +651,14 @@ def _run_remote_per_call(env, env_type: str, code: str, effective_task_id: str, env.execute(f"rm -rf {quoted_sandbox_dir}", cwd="/", timeout=15) except Exception: logger.debug("Failed to clean up remote sandbox %s", sandbox_dir) - - duration = round(time.monotonic() - exec_start, 2) - stdout_text, stdout_metadata = _clean_output(stdout_text) - result: Dict[str, Any] = { - "status": status, - "output": stdout_text, - "exit_code": exit_code, - "tool_calls_made": tool_call_counter[0], - "duration_seconds": duration, - } - result.update(stdout_metadata) - + result = _remote_result(status, stdout_text, exec_start, + {"exit_code": exit_code, "tool_calls_made": tool_call_counter[0]}) if status == "timeout": - timeout_msg = f"Script timed out after {timeout}s and was killed." - result["error"] = timeout_msg - result["output"] = _with_timeout_notice(stdout_text, timeout_msg) + _apply_timeout(result, f"Script timed out after {timeout}s and was killed.") logger.warning("execute_code (remote) timed out after %ss (limit %ss) with %d tool calls", - duration, timeout, tool_call_counter[0]) + result["duration_seconds"], timeout, tool_call_counter[0]) elif status == "interrupted": - result["output"] = _format_interrupted_output(stdout_text) + result["output"] = _format_interrupted_output(result["output"]) elif exit_code != 0: result["status"] = "error" result["error"] = f"Script exited with code {exit_code}" @@ -743,43 +667,30 @@ def _run_remote_per_call(env, env_type: str, code: str, effective_task_id: str, def _execute_remote(code: str, task_id: Optional[str], enabled_tools: Optional[List[str]], reset: bool = False) -> str: - """Run code on the remote terminal backend. - - Preferred path: the owner's persistent remote session kernel - (tools/code_kernel_remote.py — detached runner + file cell protocol). - Fallback: the per-call script ship — the fail-open route when a kernel - cannot be spawned, and the only route for hosts that cannot sustain a - background process. - """ + """Run code on the remote terminal backend: the owner's persistent remote session kernel + (tools/code_kernel_remote.py) first, else the per-call script ship — the fail-open route when + a kernel cannot be spawned and the only route for hosts that cannot sustain a background process.""" _cfg = _load_config() timeout = _cfg.get("timeout", DEFAULT_TIMEOUT) max_tool_calls = _cfg.get("max_tool_calls", DEFAULT_MAX_TOOL_CALLS) sandbox_tools = _sandbox_tools_for(enabled_tools) - effective_task_id = task_id or "default" env, env_type = _get_or_create_env(effective_task_id) exec_start = time.monotonic() - try: py_check = env.execute("command -v python3 >/dev/null 2>&1 && echo OK", cwd="/", timeout=15) if "OK" not in py_check.get("output", ""): return json.dumps({ "status": "error", - "error": ( - f"Python 3 is not available in the {env_type} terminal " - "environment. Install Python to use execute_code with " - "remote backends." - ), - "tool_calls_made": 0, - "duration_seconds": 0, + "error": (f"Python 3 is not available in the {env_type} terminal " + "environment. Install Python to use execute_code with remote backends."), + "tool_calls_made": 0, "duration_seconds": 0, }) - # Session-kernel path: one persistent kernel per owner on the # run-to-completion transport. Spawn failure falls OPEN to the per-call # path below so a degraded remote host never blocks execution. try: from tools.code_kernel_remote import execute_in_remote_kernel - kernel_result = execute_in_remote_kernel( code, env=env, env_type=env_type, task_env_id=effective_task_id, sandbox_tools=frozenset(sandbox_tools), timeout=timeout, @@ -789,13 +700,11 @@ def _execute_remote(code: str, task_id: Optional[str], enabled_tools: Optional[L except Exception: logger.warning("remote session-kernel path failed; falling back to per-call", exc_info=True) kernel_result = None - if kernel_result is not None: return _finish_remote_kernel_result(kernel_result, timeout=timeout, exec_start=exec_start) logger.info("remote session kernel unavailable on %s; using per-call path", env_type) except Exception as exc: return _remote_failure(exc, exec_start, 0) - return _run_remote_per_call(env, env_type, code, effective_task_id, sandbox_tools, timeout=timeout, max_tool_calls=max_tool_calls, exec_start=exec_start) @@ -811,47 +720,37 @@ def execute_code( enabled_tools: Optional[List[str]] = None, reset: bool = False, ) -> str: - """Run Python in the session's persistent kernel (local) or on the remote - terminal backend, with RPC access to a subset of Hermes tools; returns the - JSON result string. + """Run Python in the session's persistent kernel (local) or on the remote terminal backend, + with RPC access to a subset of Hermes tools; returns the JSON result string. - "Sandbox" here means the security envelope (env scrubbing, tool whitelist + - call budget, output redaction) — not an isolation jail: in the default - `project` mode, code runs in the session's cwd with the project venv. - ``enabled_tools`` is intersected with SANDBOX_ALLOWED_TOOLS; ``reset`` kills - the existing kernel first (ignored on per-call paths). + "Sandbox" means the security envelope (env scrubbing, tool whitelist + call budget, output + redaction), not an isolation jail: default `project` mode runs in the session's cwd with the + project venv. ``enabled_tools`` ∩ SANDBOX_ALLOWED_TOOLS; ``reset`` kills the existing kernel + first (ignored on per-call paths). """ if not SANDBOX_AVAILABLE: return tool_error( "execute_code sandbox is unavailable in this environment. " "Use normal tool calls (terminal, read_file, write_file, ...) instead." ) - - # Fail closed under a terminal-policy refusal scope: the routed profile's - # terminal policy is unresolved, so refuse rather than inherit the launch - # process's ambient policy. + # Fail closed under a terminal-policy refusal scope: the routed profile's terminal + # policy is unresolved, so refuse rather than inherit the launch process's ambient policy. try: from tools.terminal_scope import enforce_no_refusal - enforce_no_refusal() except Exception as refusal: return tool_error( f"execute_code refused: {refusal} " - "(profile terminal policy unresolved; fix the profile's " - "config.yaml / .env and retry)" + "(profile terminal policy unresolved; fix the profile's config.yaml / .env and retry)" ) - if not code or not code.strip(): return tool_error( "No code provided. execute_code requires a non-empty 'code' " - "parameter containing Python source. To run shell commands, use " - "terminal(command=...) instead." + "parameter containing Python source. To run shell commands, use terminal(command=...) instead." ) - - # Hard-block gateway-lifecycle commands (mirrors the terminal_tool guard — - # otherwise `os.system("launchctl bootout ...")` here bypasses it and - # SIGTERMs the gateway mid-task). Gated on PID-file ownership, not the - # inherited env marker. + # Hard-block gateway-lifecycle commands (mirrors the terminal_tool guard — otherwise + # `os.system("launchctl bootout ...")` here bypasses it and SIGTERMs the gateway mid-task). + # Gated on PID-file ownership, not the inherited env marker. from tools.process_registry import _is_supervised_gateway_process if _is_supervised_gateway_process(): from cron.lifecycle_guard import contains_gateway_lifecycle_command @@ -862,64 +761,47 @@ def execute_code( "it could complete (SIGTERM propagates to child processes). " "Run the lifecycle command from a shell outside the gateway." ) - from tools.terminal_tool import _get_env_config, _docker_has_host_access _env_config = _get_env_config() env_type = _env_config["env_type"] - - # Arbitrary Python never passes through terminal()/DANGEROUS_PATTERNS, so - # guard the whole script before either dispatch path spawns it — in this - # (tool-executor) thread, which holds the session context. A Docker sandbox - # with host bind mounts is not isolated, so it gets no container fast-path. + # Arbitrary Python never passes through terminal()/DANGEROUS_PATTERNS, so guard the whole + # script before either dispatch path spawns it — in this (tool-executor) thread, which holds + # the session context. A Docker sandbox with host bind mounts gets no container fast-path. from tools.approval import check_execute_code_guard _guard = check_execute_code_guard(code, env_type, has_host_access=_docker_has_host_access(_env_config)) if not _guard.get("approved", False): return _error_result(_guard.get("message") or "execute_code blocked by approval guard.") - - # Clear a stale interrupt bit that landed on this thread during the blocking - # approval-wait so it can't kill the just-approved run on the first poll - # (either dispatch path). A genuine post-clear interrupt re-sets the bit. + # Clear a stale interrupt bit that landed during the blocking approval-wait so it can't + # kill the just-approved run on the first poll. A genuine post-clear interrupt re-sets it. if _guard.get("user_approved"): from tools.interrupt import clear_current_thread_interrupt clear_current_thread_interrupt() - if env_type != "local": return _execute_remote(code, task_id, enabled_tools, reset=bool(reset)) - from tools.interrupt import is_interrupted as _is_interrupted - # Session kernels are always on locally: one interpreter per conversation; - # the guards above already ran for this cell, and the kernel path reuses the - # same env builder, RPC server, and output redaction as the remote path. + # Session kernels are always on locally (one interpreter per conversation); the guards above + # already ran for this cell, and the kernel path shares env builder, RPC server and redaction. from tools.code_kernel import execute_in_session_kernel - _cfg = _load_config() _mode = _get_execution_mode() return execute_in_session_kernel( - code, - task_id=task_id or "", - mode=_mode, - child_python=_resolve_child_python(_mode), + code, task_id=task_id or "", mode=_mode, child_python=_resolve_child_python(_mode), child_cwd=_resolve_child_cwd(_mode, "", task_id=task_id or ""), sandbox_tools=frozenset(_sandbox_tools_for(enabled_tools)), timeout=_cfg.get("timeout", DEFAULT_TIMEOUT), max_tool_calls=_cfg.get("max_tool_calls", DEFAULT_MAX_TOOL_CALLS), - reset=bool(reset), - is_interrupted=_is_interrupted, + reset=bool(reset), is_interrupted=_is_interrupted, ) def _kill_process_group(proc, escalate: bool = False): - """Kill the child and its whole process tree (cross-platform) via - agent.deadline.kill_process_tree: SIGTERM the tree (killpg + psutil - descendant sweep for setsid'd grandchildren; ``taskkill /T /F`` on Windows); - with ``escalate=True`` wait 5s then SIGKILL survivors. Never raises — a - delegation failure degrades to a plain ``proc.kill()``.""" + """Kill the child and its whole process tree via agent.deadline.kill_process_tree: SIGTERM + (killpg + psutil descendant sweep; ``taskkill /T /F`` on Windows, where sig is ignored); with + ``escalate=True`` wait 5s then SIGKILL survivors. Never raises — falls back to ``proc.kill()``.""" import signal as _signal - def _tree_signal(sig) -> None: try: from agent.deadline import kill_process_tree as _deadline_kill_tree - _deadline_kill_tree(proc.pid, sig=sig) except Exception as e: logger.debug("Could not terminate process tree: %s", e, exc_info=True) @@ -927,8 +809,6 @@ def _kill_process_group(proc, escalate: bool = False): proc.kill() except Exception as e2: logger.debug("Could not kill process: %s", e2, exc_info=True) - - # sig is ignored on Windows (taskkill /F is already forceful). _tree_signal(getattr(_signal, "SIGTERM", None)) if escalate: try: @@ -938,14 +818,10 @@ def _kill_process_group(proc, escalate: bool = False): def _load_config() -> dict: - """Load the ``code_execution`` config section via the lightweight raw reader. - - Runs while the module-level schema is built at tool discovery, so it must - not import ``cli`` (prompt_toolkit/Rich on every startup path). - """ + """``code_execution`` config section via the lightweight raw reader — runs while the + module-level schema is built at tool discovery, so it must not import ``cli``.""" try: from hermes_cli.config import read_raw_config - cfg = read_raw_config().get("code_execution", {}) return cfg if isinstance(cfg, dict) else {} except Exception: @@ -956,18 +832,15 @@ def _load_config() -> dict: # Execution mode resolution (strict vs project) # --------------------------------------------------------------------------- -# Canonical code_execution.mode values (referenced by tests and the config layer). -# Session kernels are the only local execution model; a leftover -# code_execution.kernel_mode config key is silently ignored. +# Canonical code_execution.mode values (referenced by tests and the config layer). Session +# kernels are the only local execution model; a leftover kernel_mode config key is ignored. EXECUTION_MODES = ("project", "strict") DEFAULT_EXECUTION_MODE = "project" def _get_execution_mode() -> str: - """Active execute_code mode from ``code_execution.mode`` (invalid → default - with a warning). ``project``: session cwd + active venv python so project - deps/files resolve; ``strict``: isolated temp dir + ``sys.executable``. - Env scrubbing and the tool whitelist apply identically in both.""" + """``code_execution.mode`` (invalid → default with a warning). ``project``: session cwd + active + venv python; ``strict``: isolated temp dir + ``sys.executable``. Scrubbing/whitelist apply to both.""" cfg_value = str(_load_config().get("mode", DEFAULT_EXECUTION_MODE)).strip().lower() if cfg_value in EXECUTION_MODES: return cfg_value @@ -982,8 +855,7 @@ def _get_execution_mode() -> str: # OpenAI Function-Calling Schema # --------------------------------------------------------------------------- -# Per-tool documentation lines for the execute_code description. -# Ordered to match the canonical display order. +# Per-tool documentation lines for the execute_code description, in canonical display order. _TOOL_DOC_LINES = [ ("web_search", " web_search(query: str, limit: int = 5) -> dict\n" @@ -996,8 +868,7 @@ _TOOL_DOC_LINES = [ " read_file(path: str, offset: int = 1, limit: int = 2000) -> dict\n" " Lines are 1-indexed. Returns {\"content\": \"...\", \"total_lines\": N}"), ("write_file", - " write_file(path: str, content: str) -> dict\n" - " Always overwrites the entire file."), + " write_file(path: str, content: str) -> dict\n Always overwrites the entire file."), ("search_files", " search_files(pattern: str, target=\"content\", path=\".\", file_glob=None, limit=50) -> dict\n" " target: \"content\" (search inside files) or \"files\" (find files by name). Returns {\"matches\": [...]}"), @@ -1012,22 +883,17 @@ _TOOL_DOC_LINES = [ def build_execute_code_schema(enabled_sandbox_tools: set = None, mode: str = None) -> dict: - """Build the execute_code schema listing only *enabled_sandbox_tools* — a - disabled tool (e.g. web off) must not appear or the model keeps trying it. - ``mode`` (None → current config) selects the working-directory sentence. - """ + """execute_code schema listing only *enabled_sandbox_tools* — a disabled tool (e.g. web off) + must not appear or the model keeps trying it. ``mode`` (None → config) picks the cwd sentence.""" if enabled_sandbox_tools is None: enabled_sandbox_tools = SANDBOX_ALLOWED_TOOLS if mode is None: mode = _get_execution_mode() - tool_lines = "\n".join(doc for name, doc in _TOOL_DOC_LINES if name in enabled_sandbox_tools) - import_examples = [n for n in ("web_search", "terminal") if n in enabled_sandbox_tools] if not import_examples: import_examples = sorted(enabled_sandbox_tools)[:2] import_str = ", ".join(import_examples) + ", ..." if import_examples else "..." - if mode == "strict": cwd_note = ( "Scripts run in their own temp dir, not the session's CWD — use absolute paths " @@ -1041,32 +907,26 @@ def build_execute_code_schema(enabled_sandbox_tools: set = None, "Hermes's own python (the common case — stdlib plus Hermes's " "deps; check `import x` before relying on project packages)." ) - - # Remote hosts that fail open to per-call are not worth schema words; the - # result's `kernel` field tells the truth per call. + # Remote hosts that fail open to per-call are not worth schema words; the result's + # `kernel` field tells the truth per call. description = ( "Run Python that calls Hermes tools programmatically. Use when you " "need 3+ tool calls with logic between them: filtering/reducing " "large outputs before they enter context, branching, or loops " "(N pages/files, retry on failure). Use normal tool calls for " - "single calls, results you must reason over in full, or anything " - "needing user interaction.\n\n" + "single calls, results you must reason over in full, or anything needing user interaction.\n\n" "Calls run in a persistent session kernel: variables, imports, and " "loaded data survive across execute_code calls, so build on earlier " - "work instead of re-loading it. A timed-out or interrupted call " - "loses that state.\n\n" + "work instead of re-loading it. A timed-out or interrupted call loses that state.\n\n" f"Available via `from hermes_tools import ...`:\n\n" f"{tool_lines}\n\n" "Limits: 5-minute timeout, max 50 tool calls per call. Stdout over " - "50KB shows head/tail inline; the FULL text is auto-saved to a file " - "whose path rides in the result.\n\n" + "50KB shows head/tail inline; the FULL text is auto-saved to a file whose path rides in the result.\n\n" f"{cwd_note}\n\n" "Built-in helpers (no import): json_parse(text) — tolerant " "json.loads for terminal() output; shell_quote(s) — shlex.quote for " - "dynamic shell args; retry(fn, max_attempts=3, delay=2) — " - "exponential backoff." + "dynamic shell args; retry(fn, max_attempts=3, delay=2) — exponential backoff." ) - return { "name": "execute_code", "description": description, @@ -1084,8 +944,7 @@ def build_execute_code_schema(enabled_sandbox_tools: set = None, "reset": { "type": "boolean", "description": ( - "Discard the kernel's persistent state and start " - "fresh before running this code." + "Discard the kernel's persistent state and start fresh before running this code." ), }, }, @@ -1094,14 +953,13 @@ def build_execute_code_schema(enabled_sandbox_tools: set = None, } -# Default schema used at registration time (all sandbox tools listed, -# current configured mode). model_tools.py rebuilds per-session anyway. +# Registration-time schema (all sandbox tools, configured mode); model_tools.py rebuilds per-session. EXECUTE_CODE_SCHEMA = build_execute_code_schema() def _execute_code_handler(args: dict, **kwargs) -> str: - """Redirect misdirected calls (terminal's ``command`` arg, non-string - ``code``) with an actionable error before dispatching to ``execute_code``.""" + """Redirect misdirected calls (terminal's ``command`` arg, non-string ``code``) with an + actionable error before dispatching to ``execute_code``.""" if "code" not in args and "command" in args: logger.warning("execute_code received 'command' instead of the required 'code' argument") return tool_error( @@ -1109,29 +967,18 @@ def _execute_code_handler(args: dict, **kwargs) -> str: "Python source in 'code'. Use terminal(command=...) for shell " "commands; for Python, retry as execute_code(code=...)." ) - code = args.get("code", "") if code is not None and not isinstance(code, str): return tool_error( f"execute_code received a {type(code).__name__} in 'code', but it " - "requires Python source as a string. Retry as " - "execute_code(code=\"...\")." + "requires Python source as a string. Retry as execute_code(code=\"...\")." ) - - return execute_code( - code=code or "", - task_id=kwargs.get("task_id"), - enabled_tools=kwargs.get("enabled_tools"), - reset=bool(args.get("reset", False)), - ) + return execute_code(code=code or "", task_id=kwargs.get("task_id"), + enabled_tools=kwargs.get("enabled_tools"), reset=bool(args.get("reset", False))) registry.register( - name="execute_code", - toolset="code_execution", - schema=EXECUTE_CODE_SCHEMA, - handler=_execute_code_handler, - check_fn=check_sandbox_requirements, - emoji="🐍", + name="execute_code", toolset="code_execution", schema=EXECUTE_CODE_SCHEMA, + handler=_execute_code_handler, check_fn=check_sandbox_requirements, emoji="🐍", max_result_size_chars=100_000, ) diff --git a/tools/code_kernel.py b/tools/code_kernel.py index 6a1f4a7fd1..cfc11dd9da 100644 --- a/tools/code_kernel.py +++ b/tools/code_kernel.py @@ -1,28 +1,18 @@ -"""Session-persistent Python kernels for execute_code. +"""Session-persistent Python kernels for execute_code: one child per (owner, mode, +interpreter, cwd, tool-set), one code cell per call, state survives across calls. -One Python child stays alive per (owner, mode, interpreter, cwd, tool-set) and -runs one code cell per call, so variables/imports/data survive across calls. +Constraints, in order: (1) SAME security envelope as per-call (``_build_child_env`` +scrubbing, ``_rpc_server_loop`` token + per-cell tool budget, ANSI strip + secret +redaction) — only lifetime widens. (2) A wedged kernel dies, never hangs the agent: +timeout/interrupt kills the process tree and drops the registry entry; state loss is +deliberate (a cell cannot be interrupted in place safely). (3) Env frozen at spawn: +later passthrough is invisible until ``reset=true`` (the result names the kernel). -Design constraints, in order: (1) the SAME security envelope as per-call — -``_build_child_env`` scrubbing, ``_rpc_server_loop`` with the same token and -per-cell tool budget, the same ANSI strip + secret redaction; nothing here widens -what a script can reach, only how long it lives. (2) A wedged kernel dies, never -hangs the agent: timeout or interrupt kills the whole process tree and drops the -registry entry; losing state is deliberate since one cell cannot be interrupted -in place without leaving the interpreter unknown. (3) The env is frozen at spawn: -passthrough registered later is invisible until ``reset=true`` (the result names -the kernel so this is diagnosable). - -Wire protocol (host <-> child): requests are one JSON object per stdin line -``{"id", "code"}``; responses are framed on stdout as -`` \\n`` with a per-kernel random SENTINEL from the -environment. Bytes outside frames are raw fd-level output (subprocesses inherit -the real stdout), attributed to the running cell — calls are serialized per -kernel. A script forging a frame can only fake its own cell result (same trust -position as a per-call script printing a forged success message). - -Also hosts what ``tools.code_kernel_remote`` shares: owner resolution, the -registry lifecycle, and the runner's cell-exec core. +Wire protocol: one JSON request per stdin line ``{"id", "code"}``; replies framed on +stdout as `` \\n`` with a per-kernel random SENTINEL from +the env. Bytes outside frames are raw fd output attributed to the running cell (calls +are serialized per kernel). A forged frame can only fake its own cell result. +Also hosts what ``tools.code_kernel_remote`` shares: owner resolution, registry, cell core. """ from __future__ import annotations @@ -46,13 +36,11 @@ logger = logging.getLogger(__name__) _IS_WINDOWS = sys.platform == "win32" -# Runner-side cap on captured python-level output; the host applies its own -# MAX_STDOUT truncation again. +# Runner-side cap on captured python-level output; the host re-applies its own MAX_STDOUT cap. _RUNNER_CAPTURE_BYTES = 1_000_000 -# Shared by both generated runners (which define _CAPTURE_LIMIT first): exec one -# request in the persistent GLOBALS namespace and build the response payload. -# `__name__` is `__main__` so scripts behave like the per-call path. +# Shared by both generated runners (which define _CAPTURE_LIMIT first): exec one request in the +# persistent GLOBALS namespace, build the payload. `__name__` is `__main__` as on the per-call path. RUNNER_CELL_SOURCE = '''\ GLOBALS = {"__name__": "__main__", "__builtins__": __builtins__} @@ -150,17 +138,15 @@ if __name__ == "__main__": class CellAuthority: """The approval/context identity of exactly one execute_code cell. - Interpreter state persists across cells; RPC authority must not. Each cell - installs a fresh authority — captured from the CALLING thread at cell start, - exactly what ``propagate_context_to_thread`` would capture for a per-call RPC - thread — and retires it when the cell settles, so a late tool call (a - background thread the cell left behind, a raced client write) is refused - instead of running under a stale approval/session/turn identity. + Interpreter state persists across cells; RPC authority must not. Each cell installs a + fresh authority captured from the CALLING thread at cell start (what + ``propagate_context_to_thread`` captures for a per-call RPC thread) and retires it when + the cell settles, so a late tool call (leaked background thread, raced client write) is + refused instead of running under a stale approval/session/turn identity. """ def __init__(self, task_id: str): import contextvars - self.task_id = task_id self.ctx = contextvars.copy_context() self.active = True @@ -168,12 +154,10 @@ class CellAuthority: self._callbacks = (None, None) try: from tools.thread_context import _callback_api - self._api = _callback_api() self._callbacks = (self._api[0](), self._api[1]()) except Exception: - # Fail-closed, mirroring propagate_context_to_thread: with no - # callbacks installed, dangerous approvals deny. + # Fail-closed like propagate_context_to_thread: no callbacks → dangerous approvals deny. self._api = None def retire(self) -> None: @@ -182,17 +166,13 @@ class CellAuthority: def dispatch(self, tool_name: str, tool_args: dict) -> str: """Run one tool call under THIS cell's context and callbacks.""" from tools.code_execution_tool import tool_error - if not self.active: - return tool_error( - "No active execute_code cell: the cell this kernel call " - "belonged to has settled, so its tool authority is retired." - ) + return tool_error("No active execute_code cell: the cell this kernel call " + "belonged to has settled, so its tool authority is retired.") return self.ctx.run(self._invoke, tool_name, tool_args) def _invoke(self, tool_name: str, tool_args: dict) -> str: from model_tools import handle_function_call - previous = None if self._api is not None: get_approval, get_sudo, set_approval, set_sudo = self._api @@ -259,7 +239,6 @@ class SessionKernel: self.stop_event.set() if self.alive(): from tools.code_execution_tool import _kill_process_group - _kill_process_group(self.proc, escalate=True) sock, self.server_sock = self.server_sock, None try: @@ -271,16 +250,12 @@ class SessionKernel: pass if self.tmpdir: import shutil - shutil.rmtree(self.tmpdir, ignore_errors=True) class KernelRegistry: - """Key -> kernel map plus its lock (shared with the remote registry). - - Kernels are popped under the lock and torn down outside it — teardown - may block on the child process or the remote transport. - """ + """Key -> kernel map plus its lock (shared with the remote registry). Kernels are popped + under the lock and torn down outside it — teardown may block on the child or the transport.""" def __init__(self, teardown: Callable[[Any], None]): self.kernels: Dict[Tuple, Any] = {} @@ -305,57 +280,45 @@ class KernelRegistry: _REGISTRY = KernelRegistry(lambda kernel: kernel.teardown()) _KERNELS: Dict[Tuple, SessionKernel] = _REGISTRY.kernels -# Bounded lifecycle defaults (config: code_execution.max_session_kernels / -# code_execution.kernel_idle_timeout). A long-lived gateway must never -# accumulate one live child per finished conversation: stable owner id, -# owner-teardown disposal, idle reaping, max-live bound. +# Bounded lifecycle defaults (config: code_execution.max_session_kernels / kernel_idle_timeout). +# A long-lived gateway must never accumulate one live child per finished conversation: +# stable owner id, owner-teardown disposal, idle reaping, max-live bound. DEFAULT_MAX_SESSION_KERNELS = 4 DEFAULT_KERNEL_IDLE_TIMEOUT = 1800 def _lifecycle_limits() -> Tuple[int, int]: from tools.code_execution_tool import _load_config - config = _load_config() - def limit(key: str, default: int) -> int: try: return max(1, int(config.get(key, default))) except (TypeError, ValueError): return default - return (limit("max_session_kernels", DEFAULT_MAX_SESSION_KERNELS), limit("kernel_idle_timeout", DEFAULT_KERNEL_IDLE_TIMEOUT)) def _resolve_owner(task_id: str) -> str: - """The stable identity a session kernel belongs to. + """The stable identity a session kernel belongs to: the conversation's approval session key + (context-propagated, stable across turns, distinct per session). ``run_agent`` mints a fresh + task id per turn, so a task-keyed kernel would neither survive the next turn nor be torn down + with anything; the task id is only the last-resort owner (embeds/tests without a session). - The conversation's approval session key: context-propagated, stable across - turns, distinct per session. ``run_agent`` mints a fresh task id per - top-level turn, so a task-keyed kernel would neither survive the next turn - nor ever be torn down with anything; the task id is only the last-resort - owner for embeds and tests with no session context. - - Delegated children run in a copy of the parent's context and INHERIT its - approval session key — without the ``::child::`` qualifier a child's - execute_code would attach to the parent's kernel and read its in-memory - state (verified live, both directions). Children get their own kernels, - keyed by their delegation session id. + Delegated children INHERIT the parent's approval session key — without the ``::child::`` + qualifier a child's execute_code would attach to the parent's kernel and read its state + (verified live, both directions). Children get their own kernels keyed by delegation session id. """ try: from tools.approval import get_current_session_key - session_key = get_current_session_key(default="") except Exception: session_key = "" owner = session_key or (task_id or "") try: from agent.delegation_context import is_delegated_child_context - if is_delegated_child_context(): from gateway.session_context import get_session_env - child_id = get_session_env("HERMES_SESSION_ID", "") or (task_id or "") owner = f"{owner}::child::{child_id}" except Exception: @@ -380,25 +343,16 @@ atexit.register(shutdown_all_kernels) def _rpc_forever(kernel: SessionKernel, max_tool_calls: int, sandbox_tools: frozenset) -> None: - """Serve tool RPC for the kernel's whole life. - - ``_rpc_server_loop`` serves one connection and returns on disconnect or its - 300s idle timeout; a kernel legitimately idles longer between cells, so - re-accept until teardown (the client stub reconnects: HERMES_RPC_PERSISTENT). - - The serving thread carries NO frozen authority: every dispatch routes - through the CURRENT cell's ``CellAuthority``, so a later cell's tool calls - run under that cell's context, not whatever the first cell captured. - """ + """Serve tool RPC for the kernel's whole life: ``_rpc_server_loop`` returns on disconnect or + its 300s idle timeout, and a kernel idles longer between cells, so re-accept until teardown + (the client stub reconnects: HERMES_RPC_PERSISTENT). The serving thread carries NO frozen + authority — every dispatch routes through the CURRENT cell's ``CellAuthority``.""" from tools.code_execution_tool import _rpc_server_loop, tool_error - def _dispatch(tool_name: str, tool_args: dict) -> str: authority = kernel.cell_authority if authority is None: - return tool_error("No active execute_code cell: this kernel has no cell " - "authority installed.") + return tool_error("No active execute_code cell: this kernel has no cell authority installed.") return authority.dispatch(tool_name, tool_args) - while not kernel.stop_event.is_set(): _rpc_server_loop(kernel.server_sock, "", kernel.tool_call_log, kernel.tool_call_counter, max_tool_calls, sandbox_tools, kernel.stop_event, kernel.rpc_token, @@ -408,19 +362,15 @@ def _rpc_forever(kernel: SessionKernel, max_tool_calls: int, def _stdout_reader(kernel: SessionKernel) -> None: """Split the child's stdout into protocol frames and raw passthrough.""" from tools.code_execution_tool import MAX_STDOUT_BYTES - assert kernel.proc is not None and kernel.proc.stdout is not None stream = kernel.proc.stdout marker = ("\n" + kernel.sentinel + " ").encode("utf-8") - def raw(data: bytes) -> None: kernel.raw.append(data, MAX_STDOUT_BYTES) - buf = b"" while True: - # read1: return as soon as any bytes arrive. A plain read(n) on a - # BufferedReader blocks until n bytes or EOF, which would sit on a - # complete frame smaller than the buffer forever. + # read1 returns as soon as any bytes arrive; a plain read(n) on a BufferedReader + # blocks until n bytes or EOF and would sit on a complete small frame forever. chunk = stream.read1(4096) if not chunk: if buf: @@ -431,8 +381,7 @@ def _stdout_reader(kernel: SessionKernel) -> None: while True: index = buf.find(marker) if index < 0: - # Keep a marker-sized tail in case the marker is split - # across reads; everything before it is raw output. + # Keep a marker-sized tail (marker may be split across reads); the rest is raw. spill = buf[: -len(marker)] if len(buf) > len(marker) else b"" if spill: raw(spill) @@ -448,8 +397,7 @@ def _stdout_reader(kernel: SessionKernel) -> None: try: length = int(rest[:newline]) except ValueError: - # Not a real frame header (user output that happens to - # contain the marker bytes); treat the marker as raw. + # Not a real frame header (user output containing the marker bytes): raw. raw(marker) buf = rest continue @@ -469,7 +417,6 @@ def _stdout_reader(kernel: SessionKernel) -> None: def _stderr_reader(kernel: SessionKernel) -> None: from tools.code_execution_tool import MAX_STDERR_BYTES - assert kernel.proc is not None and kernel.proc.stderr is not None while True: chunk = kernel.proc.stderr.read1(4096) @@ -501,52 +448,40 @@ def _bind_rpc_socket(kernel: SessionKernel) -> str: def _spawn(kernel: SessionKernel, *, child_python: str, child_cwd: str, sandbox_tools: frozenset, max_tool_calls: int) -> None: from tools.code_execution_tool import _build_child_env, generate_hermes_tools_module - kernel.tmpdir = tempfile.mkdtemp(prefix="hermes_kernel_") kernel.rpc_token = secrets.token_urlsafe(32) kernel.sentinel = "@@HERMES-KERNEL-" + secrets.token_urlsafe(16) + "@@" rpc_endpoint = _bind_rpc_socket(kernel) - for name, src in (("hermes_tools.py", generate_hermes_tools_module(list(sandbox_tools))), ("hermes_kernel_runner.py", KERNEL_RUNNER_SOURCE)): with open(os.path.join(kernel.tmpdir, name), "w", encoding="utf-8") as f: f.write(src) - child_env = _build_child_env(rpc_endpoint=rpc_endpoint, rpc_token=kernel.rpc_token, tmpdir=kernel.tmpdir, child_python=child_python) child_env["HERMES_KERNEL_SENTINEL"] = kernel.sentinel - # Cells clip stdout to the inline cap; the full text spills to the kernel's - # own tmpdir so the agent can read_file the middle instead of re-running. + # Full clipped stdout spills to the kernel's tmpdir so the agent can read_file the middle. child_env["HERMES_KERNEL_SPILL_DIR"] = kernel.tmpdir - # Tell the generated client to reconnect after the RPC server's idle - # timeout — a kernel outlives the 300s window between cells. + # Generated client reconnects after the RPC server's 300s idle timeout between cells. child_env["HERMES_RPC_PERSISTENT"] = "1" - kernel.proc = subprocess.Popen( [child_python, os.path.join(kernel.tmpdir, "hermes_kernel_runner.py")], - # Strict mode resolves an empty cwd: the kernel's own staging dir - # then plays the per-call tmpdir's role. + # Strict mode passes an empty cwd: the kernel's staging dir plays the per-call tmpdir's role. cwd=child_cwd or kernel.tmpdir, env=child_env, stdout=subprocess.PIPE, stderr=subprocess.PIPE, stdin=subprocess.PIPE, start_new_session=True, creationflags=subprocess.CREATE_NO_WINDOW if _IS_WINDOWS else 0, ) - - # Deliberately NOT propagate_context_to_thread: that would freeze the - # spawning cell's context/callbacks into the server thread for the - # kernel's whole life. Authority is rebound per cell via CellAuthority. + # Deliberately NOT propagate_context_to_thread: that would freeze the spawning cell's + # context/callbacks into the server thread for life. Authority is rebound per cell. for target, args in ((_rpc_forever, (kernel, max_tool_calls, sandbox_tools)), (_stdout_reader, (kernel,)), (_stderr_reader, (kernel,))): threading.Thread(target=target, args=args, daemon=True).start() def _acquire_kernel(key: Tuple, reset: bool) -> Tuple[SessionKernel, bool]: - """Look up or register the kernel for *key*; returns (kernel, state_reset). - - Every entry also sweeps idle-expired kernels and enforces the process-wide - LRU cap, so a long-lived host stays bounded even for owners that never - toggle or reset. Doomed kernels are popped under the lock, torn down outside it. - """ + """Look up or register the kernel for *key*; returns (kernel, state_reset). Every entry also + sweeps idle-expired kernels and enforces the process-wide LRU cap (doomed kernels are popped + under the lock, torn down outside it), so a long-lived host stays bounded.""" cap, idle_timeout = _lifecycle_limits() with _REGISTRY.lock: now = time.monotonic() @@ -594,23 +529,17 @@ def _cell_result(kernel: SessionKernel, key: Tuple, status: str, payload: Dict[s from tools.code_execution_tool import _sandbox_failure_hint, _truncate_stdout_text from agent.redact import redact_sensitive_text from tools.ansi_strip import strip_ansi - def clean(text: str) -> str: return redact_sensitive_text(strip_ansi(text), code_file=True) - if status in ("timeout", "interrupted"): - # No safe way to interrupt one cell in place: kill the kernel, - # report the state loss, let the next call respawn. + # No safe way to interrupt one cell in place: kill the kernel, report the loss, respawn next call. _REGISTRY.discard(key, kernel) - duration = round(time.monotonic() - exec_start, 2) kernel.execution_count = int(payload.get("execution_count", kernel.execution_count + 1)) - stderr_raw = kernel.stderr.drain() stdout_text = clean(str(payload.get("stdout", "")) + kernel.raw.drain()) cell_stderr = clean(str(payload.get("stderr", "")) + stderr_raw) stdout_text, stdout_metadata = _truncate_stdout_text(stdout_text) - cell_status = payload.get("status", "") result: Dict[str, Any] = { "status": status, "output": stdout_text, "exit_code": 0, @@ -619,9 +548,7 @@ def _cell_result(kernel: SessionKernel, key: Tuple, status: str, payload: Dict[s "execution_count": kernel.execution_count, "state_reset": state_reset}, } result.update(stdout_metadata) - - # Cell-side spill (runner clipped before replying): surface the full-output - # path with the same read_file recipe as the host-side spill. + # Cell-side spill (runner clipped before replying): same read_file recipe as the host-side spill. cell_spill = str(payload.get("stdout_spill_path", "") or "") if cell_spill and payload.get("stdout_clipped"): result["stdout_spill_path"] = cell_spill @@ -630,7 +557,6 @@ def _cell_result(kernel: SessionKernel, key: Tuple, status: str, payload: Dict[s f'— page it with read_file(path="{cell_spill}", offset=...) instead of re-running. ' "(Kernel state persists: printing a narrower slice next call is often cheaper.)" ) - if status == "timeout": message = (f"Cell timed out after {timeout}s; the session kernel was killed and its " "state was lost. The next execute_code call starts a fresh kernel.") @@ -638,7 +564,6 @@ def _cell_result(kernel: SessionKernel, key: Tuple, status: str, payload: Dict[s output=(stdout_text + "\n\n⏰ " + message) if stdout_text else ("⏰ " + message)) elif status == "interrupted": from tools.code_execution_tool import _format_interrupted_output - result.update(exit_code=-1, output=_format_interrupted_output(stdout_text), error="Interrupted; the session kernel was killed and its state was lost.") elif cell_status == "error": @@ -667,41 +592,29 @@ def execute_in_session_kernel( code: str, *, task_id: str, mode: str, child_python: str, child_cwd: str, sandbox_tools: frozenset, timeout: int, max_tool_calls: int, reset: bool, is_interrupted, ) -> str: - """Run one cell in the (owner, mode, python, cwd, tools) session kernel. - - The owner is the conversation's session key (``_resolve_owner``), not the - per-turn task id, so state survives across user turns of one conversation - and dies with the session. - """ + """Run one cell in the (owner, mode, python, cwd, tools) session kernel. The owner is the + session key (``_resolve_owner``), not the per-turn task id, so state survives across turns.""" key = (_resolve_owner(task_id) or "", mode, child_python, child_cwd, tuple(sorted(sandbox_tools))) exec_start = time.monotonic() kernel, state_reset = _acquire_kernel(key, reset) reused = kernel.proc is not None - - # Captured on the calling thread BEFORE the cell runs — the same snapshot a - # per-call RPC thread would have received — and installed atomically on the - # kernel so the serving thread dispatches this cell's tool calls under this - # cell's approval/session/turn identity. + # Captured on the calling thread BEFORE the cell runs (the snapshot a per-call RPC thread + # would get) and installed on the kernel so RPC dispatches under THIS cell's identity. authority = CellAuthority(task_id) - with kernel.lock: try: if kernel.proc is None: _spawn(kernel, child_python=child_python, child_cwd=child_cwd, sandbox_tools=sandbox_tools, max_tool_calls=max_tool_calls) assert kernel.proc is not None and kernel.proc.stdin is not None - - # Per-cell tool budget: the RPC loop enforces counter < max, so a - # fresh cell starts from zero without restarting the server. + # Per-cell tool budget: the RPC loop enforces counter < max; reset without restarting. kernel.tool_call_counter[0] = 0 # Anything raw that leaked between cells belongs to no cell. kernel.raw.drain() kernel.stderr.drain() kernel.cell_authority = authority - kernel.proc.stdin.write((json.dumps({"id": uuid.uuid4().hex, "code": code}) + "\n").encode("utf-8")) kernel.proc.stdin.flush() - status, payload = _await_cell(kernel, timeout, is_interrupted) result = _cell_result( kernel, key, status, payload, @@ -718,7 +631,6 @@ def execute_in_session_kernel( "duration_seconds": round(time.monotonic() - exec_start, 2), }, ensure_ascii=False) finally: - # The cell has settled on every path (success, exception, timeout, - # exit, kernel death): its tool authority retires with it, so + # The cell has settled on every path: its tool authority retires with it, so # nothing the cell left running can dispatch under it. authority.retire() diff --git a/tools/code_kernel_remote.py b/tools/code_kernel_remote.py index ed26e8427d..24b47830bd 100644 --- a/tools/code_kernel_remote.py +++ b/tools/code_kernel_remote.py @@ -32,15 +32,13 @@ from tools.code_kernel import RUNNER_CELL_SOURCE, KernelRegistry logger = logging.getLogger(__name__) -# How often the host polls the remote for a cell result file. Each poll is -# one env.execute round-trip (typically 0.1-0.4s on ssh/docker), so this is -# a floor, not a rate. +# Host poll interval for a cell result file; each poll is one env.execute round-trip +# (0.1-0.4s on ssh/docker), so this is a floor, not a rate. _CELL_POLL_INTERVAL = 0.5 -# The remote runner: a tiny forever-loop that polls for cell request files, -# execs them in one persistent namespace, and writes response files. It is -# deliberately transport-agnostic (pure files) and stdlib-only. Cells and -# tool-RPC share the kernel dir but use distinct prefixes. +# The remote runner: a forever-loop that polls for cell request files, execs them in one +# persistent namespace, writes response files. Pure files + stdlib only (transport-agnostic); +# cells and tool-RPC share the kernel dir under distinct prefixes. REMOTE_KERNEL_RUNNER_SOURCE = '''\ """Auto-generated Hermes REMOTE session-kernel runner (file cell protocol).""" import contextlib @@ -173,7 +171,6 @@ def _spawn_remote_kernel(env, env_type: str, owner: str, task_env_id: str, from tools.code_execution_tool import ( MAX_STDOUT_BYTES, _ship_file_to_remote, _env_temp_dir, generate_hermes_tools_module, ) - kernel_dir = f"{_env_temp_dir(env)}/hermes_rkernel_{uuid.uuid4().hex[:12]}" q_dir = shlex.quote(kernel_dir) kernel = None @@ -184,7 +181,6 @@ def _spawn_remote_kernel(env, env_type: str, owner: str, task_env_id: str, cell_source=RUNNER_CELL_SOURCE, capture_limit=MAX_STDOUT_BYTES, idle_exit=idle_exit)) _ship_file_to_remote(env, f"{kernel_dir}/hermes_tools.py", generate_hermes_tools_module(list(sandbox_tools), transport="file")) - env_prefix = ( f"HERMES_KERNEL_DIR={q_dir} " f"HERMES_RPC_DIR={shlex.quote(kernel_dir + '/rpc')} " @@ -226,10 +222,8 @@ def _acquire_remote_kernel(env, env_type: str, owner: str, task_env_id: str, """Find/respawn the owner's kernel: (kernel|None, reused, state_reset, state_lost).""" key = _kernel_key(owner, env_type, task_env_id) state_lost = state_reset = False - with _REGISTRY.lock: kernel = _REMOTE_KERNELS.get(key) - if kernel is not None and reset: _REGISTRY.discard(key, kernel) kernel, state_reset = None, True @@ -239,7 +233,6 @@ def _acquire_remote_kernel(env, env_type: str, owner: str, task_env_id: str, # best-effort dir cleanup; the process is already gone). _REGISTRY.discard(key, kernel) kernel, state_lost = None, True - reused = kernel is not None if kernel is None: kernel = _spawn_remote_kernel(env, env_type, owner, task_env_id, sandbox_tools, idle_exit=idle_exit) @@ -252,7 +245,6 @@ def _acquire_remote_kernel(env, env_type: str, owner: str, task_env_id: str, def _run_remote_cell(kernel: RemoteKernel, code: str, timeout: int) -> Tuple[str, Dict[str, Any]]: """Ship one cell request and poll for its result: (cell status, payload).""" from tools.code_execution_tool import _ship_file_to_remote - kernel.cell_seq += 1 seq = f"{kernel.cell_seq:06d}" q_cells = shlex.quote(f"{kernel.kernel_dir}/cells") @@ -260,7 +252,6 @@ def _run_remote_cell(kernel: RemoteKernel, code: str, timeout: int) -> Tuple[str _ship_file_to_remote(kernel.env, f"{kernel.kernel_dir}/cells/cell_req_{seq}.json.tmp", json.dumps({"id": seq, "code": code}, ensure_ascii=False)) kernel.sh(f"mv {q_cells}/cell_req_{seq}.json.tmp {q_cells}/cell_req_{seq}.json", timeout=10) - deadline = time.monotonic() + timeout while time.monotonic() < deadline: try: @@ -286,15 +277,13 @@ def execute_in_remote_kernel( ) -> Optional[Dict[str, Any]]: """Run one cell in the owner's remote kernel. - Returns the raw cell result dict (caller does output post-processing), - or ``None`` when no kernel could be spawned — the caller falls open to - the per-call path. ``state_lost`` / ``state_reset`` / ``reused`` ride in - the ``kernel`` sub-dict, matching the local kernel's result shape. + Returns the raw cell result dict (caller post-processes output), or ``None`` when no + kernel could be spawned (caller falls open to per-call). ``state_lost`` / ``state_reset`` / + ``reused`` ride in the ``kernel`` sub-dict, matching the local kernel's result shape. """ from tools.code_kernel import _resolve_owner from tools.code_execution_tool import _rpc_poll_loop from tools.thread_context import propagate_context_to_thread - owner = _resolve_owner(task_env_id) kernel, reused, state_reset, state_lost = _acquire_remote_kernel( env, env_type, owner, task_env_id, sandbox_tools, reset=reset, idle_exit=idle_exit) @@ -302,21 +291,17 @@ def execute_in_remote_kernel( return None # fail open to per-call key = _kernel_key(owner, env_type, task_env_id) kernel.last_used = time.monotonic() - - # Clean stale tool-RPC requests from a previous cell before arming this - # cell's poll loop, so a background thread the last cell leaked cannot - # smuggle a call into this cell's authority window. + # Clean stale tool-RPC requests from a previous cell before arming this cell's poll loop, so + # a background thread the last cell leaked cannot smuggle a call into this authority window. q_rpc = shlex.quote(kernel.kernel_dir + '/rpc') try: kernel.sh(f"rm -f {q_rpc}/req_* {q_rpc}/res_*", timeout=10) except Exception: pass - tool_call_log: list = [] tool_call_counter, stop_event = [0], threading.Event() - # Per-cell RPC thread carrying THIS call's approval/session context — - # the remote analogue of CellAuthority: authority lives exactly as long - # as the cell's poll loop. + # Per-cell RPC thread carrying THIS call's approval/session context — the remote analogue + # of CellAuthority: authority lives exactly as long as the cell's poll loop. rpc_thread = threading.Thread( target=propagate_context_to_thread(_rpc_poll_loop), args=(env, f"{kernel.kernel_dir}/rpc", task_env_id, tool_call_log, tool_call_counter, @@ -324,14 +309,12 @@ def execute_in_remote_kernel( daemon=True, ) rpc_thread.start() - cell_status, cell_payload = "no-result", {} try: cell_status, cell_payload = _run_remote_cell(kernel, code, timeout) finally: stop_event.set() rpc_thread.join(timeout=5) - kernel_info: Dict[str, Any] = {"reused": reused, "remote": True} result: Dict[str, Any] = { "status": "error", "stdout": cell_payload.get("stdout", ""), @@ -339,8 +322,7 @@ def execute_in_remote_kernel( "tool_calls_made": tool_call_counter[0], "kernel": kernel_info, } if cell_status in ("timeout", "protocol-error", "no-result"): - # No safe way to interrupt one cell in place (same contract as - # local): kill the kernel, report the loss, respawn next call. + # No safe way to interrupt one cell in place (same contract as local): kill, report, respawn. _REGISTRY.discard(key, kernel) if cell_status == "timeout": result["status"] = "timeout" @@ -350,7 +332,6 @@ def execute_in_remote_kernel( note = "Remote kernel protocol failure; kernel killed, state lost." kernel_info.update(ended=True, state_lost=True, note=note) return result - if cell_status == "exit": _REGISTRY.discard(key, kernel) kernel_info["ended"] = True @@ -365,8 +346,7 @@ def execute_in_remote_kernel( if state_lost: kernel_info.update(state_lost=True, note=( "The previous remote kernel was gone (transport drop, container " - "restart, or idle self-exit); state from earlier calls was lost " - "and a fresh kernel was started.")) + "restart, or idle self-exit); state from earlier calls was lost and a fresh kernel was started.")) if cell_status == "error" and result["traceback"]: result["error"] = result["traceback"].strip().splitlines()[-1] return result