From 1a748cc85a43a6c1dc5c9a391706b7d607dec24b Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Sat, 26 Sep 2026 22:40:07 +0530 Subject: [PATCH] refactor(code-execution): one checked-execute helper, complete launch command Gate follow-ups on the remote lockdown: - _execute_checked(env, cmd, what, **kw) in code_execution_rpc replaces the three copy-pasted "execute, raise if returncode != 0" blocks (per-call setup, kernel dir setup, file ship via _remote_write(check=True)). env.execute always returns a dict, so the isinstance/(r or {}) guards go; the error carries command output only, never the payload. - _ship_env_file_and_launch_prefix returned a half-built "( ... && " that both callers had to close; a caller that dropped the ")" or composed it differently would lose the load-bearing subshell. It now takes the launch command, builds the shared env map (RPC dir, token, PYTHONDONTWRITEBYTECODE, routed TZ) itself, and returns the complete command; the kernel passes only HERMES_KERNEL_DIR/PYTHONPATH, which drops its duplicated TZ block and lazy hermes_time import. - _private_dirs_cmd(root, *subdirs): every caller spelled each path twice for mkdir and chmod. - _run_remote_cell publishes the cell request with one atomic _remote_write instead of ship-to-.tmp then a separate unchecked mv, saving a backend round-trip per cell. --- tools/code_execution_rpc.py | 39 ++++++++++++++++------- tools/code_execution_tool.py | 61 +++++++++++++++++------------------- tools/code_kernel_remote.py | 42 +++++++++---------------- tools/tool_result_storage.py | 2 +- 4 files changed, 72 insertions(+), 72 deletions(-) diff --git a/tools/code_execution_rpc.py b/tools/code_execution_rpc.py index 46e777f387..e1a5400149 100644 --- a/tools/code_execution_rpc.py +++ b/tools/code_execution_rpc.py @@ -30,19 +30,33 @@ def _default_dispatch(task_id): return lambda tool_name, tool_args: handle_function_call(tool_name, tool_args, task_id=task_id) -def _private_dirs_cmd(mkdir_dirs, chmod_dirs) -> str: - """Shell command creating remote dirs owner-only on a shared host. ``umask 077`` - makes intermediates and leaves private at creation (no mkdir-then-chmod window - where a co-tenant could open a dir fd); ``chmod`` then repairs a dir that - already existed with permissive modes.""" - mkdir = " ".join(shlex.quote(d) for d in mkdir_dirs) - chmod = " ".join(shlex.quote(d) for d in chmod_dirs) +def _private_dirs_cmd(root: str, *subdirs: str) -> str: + """Shell command creating *root* (and optional *subdirs* under it) owner-only + on a shared host. ``umask 077`` makes intermediates and leaves private at + creation (no mkdir-then-chmod window where a co-tenant could open a dir fd); + ``chmod`` then repairs any named dir that already existed with permissive + modes.""" + mkdir = " ".join(shlex.quote(d) for d in (subdirs or (root,))) + chmod = " ".join(shlex.quote(d) for d in (root, *subdirs)) return f"umask 077 && mkdir -p {mkdir} && chmod 700 {chmod}" +def _execute_checked(env, cmd: str, what: str, *, timeout: int, **kwargs) -> dict: + """Run *cmd* from ``/`` and raise ``RuntimeError`` on a non-zero exit. + + Used where a silent failure would ship secrets or code into a missing, + half-written, or still-permissive remote path. The error carries the + command output only, never the payload.""" + result = env.execute(cmd, cwd="/", timeout=timeout, **kwargs) + if result.get("returncode", 1) != 0: + raise RuntimeError(f"{what} failed: {result.get('output')!r}") + return result + + def _remote_write(env, remote_path: str, content: str, *, atomic: bool = False, - timeout: int = 30): - """Write *content* owner-only to *remote_path*; returns the execute() result. + timeout: int = 30, check: bool = False): + """Write *content* owner-only to *remote_path*; returns the execute() result + (``check=True`` raises on failure via _execute_checked). The base64 payload always travels as ``stdin_data``: pipe-mode backends (ssh, docker, local, singularity: the real shared-host ones) deliver it on @@ -55,8 +69,11 @@ def _remote_write(env, remote_path: str, content: str, *, atomic: bool = False, target = shlex.quote(remote_path) write = (f"base64 -d > {target}.tmp && mv -f {target}.tmp {target}" if atomic else f"base64 -d > {target}") - return env.execute(f"umask 077 && {write}", cwd="/", timeout=timeout, - stdin_data=encoded) + cmd = f"umask 077 && {write}" + if check: + return _execute_checked(env, cmd, f"remote file ship for {remote_path!r}", + timeout=timeout, stdin_data=encoded) + return env.execute(cmd, cwd="/", timeout=timeout, stdin_data=encoded) def _rpc_token_ok(request: dict, rpc_token: str) -> bool: diff --git a/tools/code_execution_tool.py b/tools/code_execution_tool.py index 6a9470626a..eeb8a475df 100644 --- a/tools/code_execution_tool.py +++ b/tools/code_execution_tool.py @@ -29,7 +29,9 @@ from tools.registry import registry, tool_error from hermes_time import get_timezone_name from tools.code_execution_env import _resolve_child_cwd, _resolve_child_python -from tools.code_execution_rpc import _private_dirs_cmd, _remote_write, _rpc_poll_loop +from tools.code_execution_rpc import ( + _execute_checked, _private_dirs_cmd, _remote_write, _rpc_poll_loop, +) from tools.tool_output_truncate import head_tail_split, truncation_notice logger = logging.getLogger(__name__) @@ -452,24 +454,21 @@ def _get_or_create_env(task_id: str): return env, env_type -def _ship_file_to_remote(env, remote_path: str, content: str) -> None: +def _ship_file_to_remote(env, remote_path: str, content: str, *, atomic: bool = False) -> None: """Write *content* owner-only to *remote_path*; the payload rides - ``stdin_data``, never an ``echo`` argv (see _remote_write). Raises on write failure: the - caller ships secrets and code into dirs it believes are locked down, so a - silent failure would run the next step against a missing or half-written - file.""" - result = _remote_write(env, remote_path, content) - if not isinstance(result, dict) or result.get("returncode", 1) != 0: - raise RuntimeError( - f"remote file ship failed for {remote_path!r}: " - f"{(result or {}).get('output', result)!r}") + ``stdin_data``, never an ``echo`` argv (see _remote_write). Raises on write + failure: the caller ships secrets and code into dirs it believes are locked + down, so a silent failure would run the next step against a missing or + half-written file.""" + _remote_write(env, remote_path, content, atomic=atomic, check=True) -def _ship_env_file_and_launch_prefix(env, remote_dir: str, env_name: str, - env_map: dict) -> str: - """Ship *env_map* as KEY=value lines to ``remote_dir/env_name`` and return a - command prefix that sources it inside a subshell; the caller appends the - launch command and the closing ``)``. +def _ship_env_file_and_launch(env, remote_dir: str, env_name: str, launch: str, *, + rpc_dir: str, rpc_token: str, **extra_env: str) -> str: + """Ship the sandbox env (RPC dir + token, PYTHONDONTWRITEBYTECODE, the routed + profile's TZ, plus *extra_env*) as KEY=value lines to ``remote_dir/env_name`` + and return the complete command that sources it inside a subshell and runs + *launch* there. The subshell is load-bearing: every env.execute() runs inside the backend's session wrapper, which re-dumps ``export -p`` into the shared session @@ -478,9 +477,15 @@ def _ship_env_file_and_launch_prefix(env, remote_dir: str, env_name: str, them into every later command on the backend (the #71296 snapshot-leak class). The token also stays off the remote shell's argv, which co-tenant users can read via ps for the command's lifetime.""" + env_map = {"HERMES_RPC_DIR": rpc_dir, "HERMES_RPC_TOKEN": rpc_token, + "PYTHONDONTWRITEBYTECODE": "1", **extra_env} + tz = get_timezone_name() # routed profile's timezone, not the bridged default's + if tz: + env_map["TZ"] = tz lines = "".join(f"{k}={shlex.quote(v)}\n" for k, v in env_map.items()) _ship_file_to_remote(env, f"{remote_dir}/{env_name}", lines) - return f"cd {shlex.quote(remote_dir)} && ( set -a && . ./{env_name} && set +a && " + return (f"cd {shlex.quote(remote_dir)} && " + f"( set -a && . ./{env_name} && set +a && {launch} )") def _env_temp_dir(env: Any) -> str: @@ -596,12 +601,8 @@ def _run_remote_per_call(env, env_type: str, code: str, effective_task_id: str, # Private dirs: the sandbox lives under a shared temp dir and carries the # RPC token (in req files) and tool results. Fail closed on setup # failure rather than ship secrets into a dir that stayed permissive. - setup = env.execute( - _private_dirs_cmd([f"{sandbox_dir}/rpc"], [sandbox_dir, f"{sandbox_dir}/rpc"]), - cwd="/", timeout=10) - if not isinstance(setup, dict) or setup.get("returncode", 1) != 0: - raise RuntimeError( - f"remote sandbox setup failed: {(setup or {}).get('output', setup)!r}") + _execute_checked(env, _private_dirs_cmd(sandbox_dir, f"{sandbox_dir}/rpc"), + "remote sandbox setup", timeout=10) rpc_token = secrets.token_urlsafe(32) _ship_file_to_remote(env, f"{sandbox_dir}/hermes_tools.py", generate_hermes_tools_module(list(sandbox_tools), transport="file")) @@ -617,17 +618,11 @@ def _run_remote_per_call(env, env_type: str, code: str, effective_task_id: str, # The token travels in a sourced env file, never in argv. No umask on # the launch command: the 700 dirs + explicit 0600 writes cover Hermes' # files, and user code keeps the remote's default file modes. - env_map = {"HERMES_RPC_DIR": f"{sandbox_dir}/rpc", - "HERMES_RPC_TOKEN": rpc_token, - "PYTHONDONTWRITEBYTECODE": "1"} - tz = get_timezone_name() # routed profile's timezone, not the bridged default's - if tz: - env_map["TZ"] = tz - launch_prefix = _ship_env_file_and_launch_prefix( - env, sandbox_dir, "sandbox.env", env_map) + launch_cmd = _ship_env_file_and_launch( + env, sandbox_dir, "sandbox.env", "exec python3 script.py", + rpc_dir=f"{sandbox_dir}/rpc", rpc_token=rpc_token) logger.info("Executing code on %s backend (task %s)...", env_type, effective_task_id[:8]) - script_result = env.execute(f"{launch_prefix} exec python3 script.py )", - timeout=timeout) + script_result = env.execute(launch_cmd, timeout=timeout) stdout_text = script_result.get("output", "") or "" exit_code = script_result.get("returncode", -1) # Backend exit codes: 124 = timeout wrapper, 130 = SIGINT. diff --git a/tools/code_kernel_remote.py b/tools/code_kernel_remote.py index 3564754501..709536309d 100644 --- a/tools/code_kernel_remote.py +++ b/tools/code_kernel_remote.py @@ -207,11 +207,10 @@ atexit.register(shutdown_all_remote_kernels) def _spawn_remote_kernel(env, env_type: str, owner: str, task_env_id: str, sandbox_tools: frozenset, *, idle_exit: int) -> Optional[RemoteKernel]: """Start a detached kernel runner on the remote. None on failure (dir removed).""" - from hermes_time import get_timezone_name - from tools.code_execution_rpc import _private_dirs_cmd + from tools.code_execution_rpc import _execute_checked, _private_dirs_cmd from tools.code_execution_tool import ( MAX_STDOUT_BYTES, _ship_file_to_remote, _env_temp_dir, - _ship_env_file_and_launch_prefix, generate_hermes_tools_module, + _ship_env_file_and_launch, generate_hermes_tools_module, ) kernel_dir = f"{_env_temp_dir(env)}/hermes_rkernel_{uuid.uuid4().hex[:12]}" q_dir = shlex.quote(kernel_dir) @@ -221,37 +220,26 @@ def _spawn_remote_kernel(env, env_type: str, owner: str, task_env_id: str, # the RPC token (in req files), tool results, and cell code/output. # Fail closed on setup failure rather than ship secrets into a dir that # stayed permissive. - setup = env.execute( - _private_dirs_cmd([f"{kernel_dir}/cells", f"{kernel_dir}/rpc"], - [kernel_dir, f"{kernel_dir}/cells", f"{kernel_dir}/rpc"]), - cwd="/", timeout=15) - if not isinstance(setup, dict) or setup.get("returncode", 1) != 0: - raise RuntimeError( - f"remote kernel dir setup failed: {(setup or {}).get('output', setup)!r}") + _execute_checked(env, _private_dirs_cmd(kernel_dir, f"{kernel_dir}/cells", + f"{kernel_dir}/rpc"), + "remote kernel dir setup", timeout=15) rpc_token = secrets.token_urlsafe(32) _ship_file_to_remote(env, f"{kernel_dir}/kernel_runner.py", REMOTE_KERNEL_RUNNER_SOURCE.format( 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_map = {"HERMES_KERNEL_DIR": kernel_dir, - "HERMES_RPC_DIR": f"{kernel_dir}/rpc", - "HERMES_RPC_TOKEN": rpc_token, - "PYTHONDONTWRITEBYTECODE": "1", - "PYTHONPATH": kernel_dir} - tz = get_timezone_name() # routed profile's timezone, matching the per-call path - if tz: - env_map["TZ"] = tz - launch_prefix = _ship_env_file_and_launch_prefix( - env, kernel_dir, "kernel.env", env_map) # kernel.env is removed after sourcing: the runner's env keeps the # values, so the token file need not sit at rest for the kernel's # lifetime. runner.log is pre-created 600 so the launch redirect never # lands at the remote's default umask. The inner `&` stays inside the # subshell where `$!` resolves to the runner pid. - started = _sh(env, f"{launch_prefix} rm -f ./kernel.env && " - f"touch runner.log && chmod 600 runner.log && " - f"{{ nohup python3 kernel_runner.py > runner.log 2>&1 & " - f'echo "PID:$!"; }} )', timeout=20) + launch_cmd = _ship_env_file_and_launch( + env, kernel_dir, "kernel.env", + "rm -f ./kernel.env && touch runner.log && chmod 600 runner.log && " + '{ nohup python3 kernel_runner.py > runner.log 2>&1 & echo "PID:$!"; }', + rpc_dir=f"{kernel_dir}/rpc", rpc_token=rpc_token, + HERMES_KERNEL_DIR=kernel_dir, PYTHONPATH=kernel_dir) + started = _sh(env, launch_cmd, timeout=20) pid = next((line.strip()[4:].strip() for line in started.splitlines() if line.strip().startswith("PID:")), "") if not pid.isdigit(): @@ -317,9 +305,9 @@ def _run_remote_cell(kernel: RemoteKernel, code: str, timeout: int) -> Tuple[str kernel.cell_seq += 1 seq = f"{kernel.cell_seq:06d}" q_cells, q_res = shlex.quote(f"{kernel.kernel_dir}/cells"), shlex.quote(f"cell_res_{seq}.json") - _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) + # One round-trip: tmp write + rename publishes the request atomically. + _ship_file_to_remote(kernel.env, f"{kernel.kernel_dir}/cells/cell_req_{seq}.json", + json.dumps({"id": seq, "code": code}, ensure_ascii=False), atomic=True) deadline = time.monotonic() + timeout while time.monotonic() < deadline: try: diff --git a/tools/tool_result_storage.py b/tools/tool_result_storage.py index e82c010dec..958d22a7bd 100644 --- a/tools/tool_result_storage.py +++ b/tools/tool_result_storage.py @@ -186,7 +186,7 @@ def _write_to_sandbox(content: str, remote_path: str, env) -> bool: # Private dir: archived results carry tool output (can hold secrets) under a # shared temp root on remote backends. The umask also covers the cat redirect. from tools.code_execution_rpc import _private_dirs_cmd - cmd = (f"{_private_dirs_cmd([storage_dir], [storage_dir])} " + cmd = (f"{_private_dirs_cmd(storage_dir)} " f"&& cat > {shlex.quote(remote_path)}") if env.execute(cmd, timeout=30, stdin_data=content).get("returncode", 1) != 0: return False