refactor(tools): compact environment base/docker/singularity modules, extract docker __init__ phases
This commit is contained in:
@@ -1,13 +1,6 @@
|
||||
"""Hermes execution environment backends.
|
||||
|
||||
Each backend provides the same interface (BaseEnvironment ABC) for running
|
||||
shell commands in a specific execution context: local, Docker, SSH,
|
||||
Singularity, Modal, Daytona, or Vercel Sandbox. (Modal additionally has
|
||||
direct and Nous-managed modes, selected via terminal.modal_mode.)
|
||||
|
||||
The terminal_tool.py factory (_create_environment) selects the backend
|
||||
based on the TERMINAL_ENV configuration.
|
||||
"""
|
||||
"""Hermes execution environment backends: one BaseEnvironment ABC for running shell commands
|
||||
in a specific context (local, Docker, SSH, Singularity, Modal direct/Nous-managed, Daytona,
|
||||
Vercel Sandbox). ``terminal_tool._create_environment`` selects the backend from TERMINAL_ENV."""
|
||||
|
||||
from tools.environments.base import BaseEnvironment
|
||||
|
||||
|
||||
+96
-202
@@ -3,12 +3,9 @@
|
||||
Unified spawn-per-call model: every command spawns a fresh ``bash -c`` process.
|
||||
A session snapshot (env vars, functions, aliases) is captured once at init and
|
||||
re-sourced before each command. CWD persists via in-band stdout markers (remote)
|
||||
or a temp file (local).
|
||||
|
||||
Cohesive pieces live in sibling modules and are re-exported here so
|
||||
``from tools.environments.base import X`` / ``patch("tools.environments.base.X")``
|
||||
keep working: ``base_output`` (collector, ProcessHandle, stdin/drain plumbing),
|
||||
``base_session_env`` (snapshot/wrapper shell scripting), ``base_wait`` (tracing).
|
||||
or a temp file (local). Cohesive pieces live in sibling modules (``base_output``,
|
||||
``base_session_env``, ``base_wait``, ``path_utils``) and are re-exported here so
|
||||
``from tools.environments.base import X`` / ``patch("tools.environments.base.X")`` keep working.
|
||||
"""
|
||||
|
||||
import json
|
||||
@@ -33,8 +30,7 @@ from tools.environments.base_output import ( # noqa: F401
|
||||
_new_output_collector,
|
||||
_pipe_stdin,
|
||||
_popen_bash,
|
||||
_start_drain_thread,
|
||||
)
|
||||
_start_drain_thread)
|
||||
from tools.environments.base_session_env import ( # noqa: F401
|
||||
_SHELL_ENV_NAME_RE,
|
||||
_SNAP_TMP,
|
||||
@@ -44,15 +40,13 @@ from tools.environments.base_session_env import ( # noqa: F401
|
||||
_export_dump_excluding_session_vars,
|
||||
_snapshot_bootstrap_script,
|
||||
_split_cwd_marker,
|
||||
_wrap_command_script,
|
||||
)
|
||||
_wrap_command_script)
|
||||
from tools.environments.base_wait import _WaitTrace
|
||||
from tools.environments.path_utils import ( # noqa: F401
|
||||
_SANDBOX_DIR_HASH_LEN,
|
||||
_SANDBOX_DIR_MAX_LEN,
|
||||
_SANDBOX_DIR_UNSAFE_RE,
|
||||
sanitize_task_id_for_path,
|
||||
)
|
||||
sanitize_task_id_for_path)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -60,11 +54,9 @@ logger = logging.getLogger(__name__)
|
||||
# (HERMES_DEBUG_INTERRUPT=1). Off by default to avoid flooding gateway logs.
|
||||
_DEBUG_INTERRUPT = bool(os.getenv("HERMES_DEBUG_INTERRUPT"))
|
||||
|
||||
# Extra seconds the ``run_bounded_sync`` backstop waits past the inner
|
||||
# ``_wait_for_process`` deadline. The inner poll loop returns partial output +
|
||||
# returncode 124; the outer bound only fires when that loop itself never
|
||||
# returns (a blocked wait that silently disables asyncio timers). Keep it
|
||||
# small so a healthy timeout still comes from the inner path.
|
||||
# Extra seconds the ``run_bounded_sync`` backstop waits past the inner ``_wait_for_process``
|
||||
# deadline: the inner loop returns partial output + 124; the outer bound only fires when that
|
||||
# loop never returns. Keep small so a healthy timeout still comes from the inner path.
|
||||
_EXECUTE_WAIT_BOUND_GRACE_S = 2.0
|
||||
|
||||
if _DEBUG_INTERRUPT:
|
||||
@@ -78,15 +70,11 @@ _activity_callback_local = threading.local()
|
||||
|
||||
|
||||
class EnvironmentConnectionError(RuntimeError):
|
||||
"""Infrastructure/connection-class failure of a terminal backend.
|
||||
|
||||
Raised when the backend itself is unreachable (SSH host down, Docker daemon
|
||||
not running, remote file sync failing on a dead link) — never for a command
|
||||
that merely exited nonzero. Subclassing RuntimeError keeps every existing
|
||||
``except RuntimeError`` catcher working. ``terminal_tool`` turns this into a
|
||||
structured ``status: "degraded"`` result (config ``terminal.degraded_mode``);
|
||||
the failed backend is never cached, so a later call retries from scratch.
|
||||
"""
|
||||
"""Infrastructure/connection-class failure of a terminal backend (SSH host down, Docker
|
||||
daemon not running, remote sync on a dead link) — never a command that merely exited
|
||||
nonzero. Subclassing RuntimeError keeps every ``except RuntimeError`` catcher working.
|
||||
``terminal_tool`` turns this into a structured ``status: "degraded"`` result; the failed
|
||||
backend is never cached, so a later call retries from scratch."""
|
||||
|
||||
def __init__(self, reason: str, *, retry_hint: str = ""):
|
||||
super().__init__(reason)
|
||||
@@ -95,8 +83,7 @@ class EnvironmentConnectionError(RuntimeError):
|
||||
"This is an infrastructure failure, not a command failure. "
|
||||
"Verify the backend is reachable (network, service running, "
|
||||
"credentials), then retry the same command — recovery is "
|
||||
"automatic once the backend is back."
|
||||
)
|
||||
"automatic once the backend is back.")
|
||||
|
||||
|
||||
def set_activity_callback(cb: Callable[[str], None] | None) -> None:
|
||||
@@ -105,31 +92,21 @@ def set_activity_callback(cb: Callable[[str], None] | None) -> None:
|
||||
|
||||
|
||||
def get_activity_callback() -> Callable[[str], None] | None:
|
||||
"""Return the thread-local activity callback (see ``set_activity_callback``).
|
||||
|
||||
For callers that must capture the calling thread's callback before handing
|
||||
work to another thread — a freshly spawned thread cannot read it back.
|
||||
"""
|
||||
"""Thread-local activity callback; capture it before handing work to another thread."""
|
||||
return getattr(_activity_callback_local, "callback", None)
|
||||
|
||||
|
||||
def touch_activity_if_due(state: dict, label: str) -> None:
|
||||
"""Fire the activity callback at most once every ``state['interval']`` seconds.
|
||||
|
||||
*state* must contain ``last_touch`` and ``start`` (monotonic timestamps);
|
||||
optional ``interval`` overrides the default 10 s cadence. Swallows all
|
||||
exceptions so callers don't need their own try/except.
|
||||
"""
|
||||
"""Fire the activity callback at most once every ``state['interval']`` (default 10 s).
|
||||
*state* holds ``last_touch``/``start`` monotonic timestamps. Swallows all exceptions."""
|
||||
now = time.monotonic()
|
||||
interval = state.get("interval", 10.0)
|
||||
if now - state["last_touch"] < interval:
|
||||
if now - state["last_touch"] < state.get("interval", 10.0):
|
||||
return
|
||||
state["last_touch"] = now
|
||||
try:
|
||||
cb = get_activity_callback()
|
||||
if cb:
|
||||
elapsed = int(now - state["start"])
|
||||
cb(f"{label} ({elapsed}s elapsed)")
|
||||
cb(f"{label} ({int(now - state['start'])}s elapsed)")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
@@ -145,12 +122,10 @@ def get_sandbox_dir() -> Path:
|
||||
|
||||
def _load_json_store(path: Path) -> dict:
|
||||
"""Load a JSON file as a dict, returning ``{}`` on any error."""
|
||||
if path.exists():
|
||||
try:
|
||||
return json.loads(path.read_text(encoding="utf-8"))
|
||||
except Exception:
|
||||
pass
|
||||
return {}
|
||||
try:
|
||||
return json.loads(path.read_text(encoding="utf-8"))
|
||||
except Exception:
|
||||
return {}
|
||||
|
||||
|
||||
def _save_json_store(path: Path, data: dict) -> None:
|
||||
@@ -168,25 +143,16 @@ def _file_mtime_key(host_path: str) -> tuple[float, int] | None:
|
||||
return None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# BaseEnvironment
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class BaseEnvironment(ABC):
|
||||
"""Common interface and unified execution flow for all Hermes backends.
|
||||
|
||||
Subclasses implement ``_run_bash()`` and ``cleanup()``. The base class
|
||||
provides ``execute()`` with session snapshot sourcing, CWD tracking,
|
||||
interrupt handling, and timeout enforcement.
|
||||
"""
|
||||
"""Common interface and unified execution flow for all Hermes backends. Subclasses
|
||||
implement ``_run_bash()`` and ``cleanup()``; the base provides ``execute()`` with
|
||||
snapshot sourcing, CWD tracking, interrupt handling and timeout enforcement."""
|
||||
|
||||
# Subclasses that embed stdin as a heredoc (Modal, Daytona) set this.
|
||||
_stdin_mode: str = "pipe" # "pipe" or "heredoc"
|
||||
|
||||
# True only when commands execute on the SAME host as the Hermes process
|
||||
# (LocalEnvironment). Controller-host facts (sys.platform, Path.home())
|
||||
# describe the execution target only when this is True.
|
||||
# (LocalEnvironment); controller-host facts then describe the execution target.
|
||||
is_local: bool = False
|
||||
|
||||
# Snapshot creation timeout (override for slow cold-starts).
|
||||
@@ -222,12 +188,7 @@ class BaseEnvironment(ABC):
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
def _run_bash(
|
||||
self,
|
||||
cmd_string: str,
|
||||
*,
|
||||
login: bool = False,
|
||||
timeout: int = 120,
|
||||
stdin_data: str | None = None,
|
||||
self, cmd_string: str, *, login: bool = False, timeout: int = 120, stdin_data: str | None = None,
|
||||
) -> ProcessHandle:
|
||||
"""Spawn a bash process to run *cmd_string*; every backend overrides this."""
|
||||
raise NotImplementedError(f"{type(self).__name__} must implement _run_bash()")
|
||||
@@ -246,66 +207,47 @@ class BaseEnvironment(ABC):
|
||||
return ()
|
||||
|
||||
def _snapshot_excluded_passthrough_names(self) -> tuple[str, ...]:
|
||||
"""Profile-scoped names that must not persist in the snapshot.
|
||||
|
||||
Monotonic for the environment lifetime: an allowlist can be cleared
|
||||
after a value was captured, and retaining the exclusion keeps that old
|
||||
value from leaking to a later profile through the shared snapshot.
|
||||
"""
|
||||
"""Profile-scoped names that must not persist in the snapshot. Monotonic for the
|
||||
environment lifetime: an allowlist can be cleared after a value was captured, and
|
||||
retaining the exclusion keeps that old value from leaking to a later profile."""
|
||||
if not self._profile_scoped_passthrough:
|
||||
return ()
|
||||
try:
|
||||
from agent.secret_scope import is_multiplex_active
|
||||
if is_multiplex_active():
|
||||
from tools.env_passthrough import get_all_passthrough
|
||||
names = (
|
||||
*get_all_passthrough(),
|
||||
*self._additional_profile_scoped_passthrough_names(),
|
||||
)
|
||||
names = (*get_all_passthrough(), *self._additional_profile_scoped_passthrough_names())
|
||||
self._snapshot_passthrough_names.update(
|
||||
name
|
||||
for name in names
|
||||
if isinstance(name, str) and _SHELL_ENV_NAME_RE.fullmatch(name)
|
||||
)
|
||||
name for name in names if isinstance(name, str) and _SHELL_ENV_NAME_RE.fullmatch(name))
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"Could not refresh profile-scoped snapshot exclusions",
|
||||
exc_info=True,
|
||||
)
|
||||
logger.debug("Could not refresh profile-scoped snapshot exclusions", exc_info=True)
|
||||
return tuple(sorted(self._snapshot_passthrough_names))
|
||||
|
||||
def init_session(self):
|
||||
"""Capture the login shell environment into the snapshot file.
|
||||
|
||||
Called once after construction. On success ``_snapshot_ready`` is set so
|
||||
commands source the snapshot instead of running under ``bash -l``. On
|
||||
failure, fall back to ``bash -l`` per command — unless a non-login probe
|
||||
shows login bash itself is dead, in which case prefer ``bash -c``.
|
||||
"""
|
||||
bootstrap = _snapshot_bootstrap_script(
|
||||
# ``_quote_cwd_for_cd`` / ``_quote_shell_path`` (not bare shlex.quote)
|
||||
# let the Windows subclass rewrite ``C:\\...`` to ``/c/...`` so the
|
||||
# bootstrap ``cd`` resolves and MSYS doesn't choke on drive paths.
|
||||
quoted_cwd=self._quote_cwd_for_cd(self.cwd),
|
||||
def _snapshot_script_kwargs(self, cwd: str) -> dict:
|
||||
"""Quoting inputs shared by the bootstrap and per-command wrapper scripts.
|
||||
``_quote_cwd_for_cd`` / ``_quote_shell_path`` (not bare shlex.quote) let the Windows
|
||||
subclass rewrite ``C:\\...`` to ``/c/...`` so ``cd`` resolves and MSYS doesn't choke."""
|
||||
return dict(
|
||||
quoted_cwd=self._quote_cwd_for_cd(cwd),
|
||||
quoted_snap=self._quote_shell_path(self._snapshot_path),
|
||||
snap_tmp_template=self._quote_shell_path(self._snapshot_path + _SNAP_TMP_SUFFIX),
|
||||
excluded_names=self._snapshot_excluded_passthrough_names(),
|
||||
cwd_marker=self._cwd_marker,
|
||||
)
|
||||
cwd_marker=self._cwd_marker)
|
||||
|
||||
def init_session(self):
|
||||
"""Capture the login shell environment into the snapshot file (once, after construction).
|
||||
On success ``_snapshot_ready`` is set so commands source the snapshot instead of running
|
||||
under ``bash -l``. On failure, fall back to ``bash -l`` per command — unless a non-login
|
||||
probe shows login bash itself is dead, in which case prefer ``bash -c``."""
|
||||
bootstrap = _snapshot_bootstrap_script(
|
||||
excluded_names=self._snapshot_excluded_passthrough_names(), **self._snapshot_script_kwargs(self.cwd))
|
||||
try:
|
||||
proc = self._run_bash(bootstrap, login=True, timeout=self._snapshot_timeout)
|
||||
result = self._wait_for_process(proc, timeout=self._snapshot_timeout)
|
||||
if int(result.get("returncode") or 0) != 0:
|
||||
raise RuntimeError(
|
||||
f"snapshot bootstrap failed with exit code {result.get('returncode')}"
|
||||
)
|
||||
raise RuntimeError(f"snapshot bootstrap failed with exit code {result.get('returncode')}")
|
||||
self._snapshot_ready = True
|
||||
self._update_cwd(result)
|
||||
logger.info(
|
||||
"Session snapshot created (session=%s, cwd=%s)",
|
||||
self._session_id,
|
||||
self.cwd,
|
||||
)
|
||||
logger.info("Session snapshot created (session=%s, cwd=%s)", self._session_id, self.cwd)
|
||||
except Exception as exc:
|
||||
self._snapshot_ready = False
|
||||
self._prefer_nonlogin, detail = self._probe_nonlogin_fallback(str(exc))
|
||||
@@ -313,16 +255,12 @@ class BaseEnvironment(ABC):
|
||||
logger.warning(
|
||||
"init_session failed (session=%s): %s — "
|
||||
"login bash unusable; falling back to non-login bash -c",
|
||||
self._session_id,
|
||||
exc,
|
||||
)
|
||||
self._session_id, exc)
|
||||
else:
|
||||
logger.warning(
|
||||
"init_session failed (session=%s): %s — "
|
||||
"falling back to bash -l per command",
|
||||
self._session_id,
|
||||
detail,
|
||||
)
|
||||
self._session_id, detail)
|
||||
|
||||
def _probe_nonlogin_fallback(self, detail: str) -> tuple[bool, str]:
|
||||
"""Run ``true`` under non-login bash; return ``(prefer_nonlogin, detail)``."""
|
||||
@@ -359,17 +297,12 @@ class BaseEnvironment(ABC):
|
||||
return shlex.quote(path)
|
||||
|
||||
def _wrap_command(self, command: str, cwd: str) -> str:
|
||||
"""Build the full bash script that sources snapshot, cd's, runs command,
|
||||
re-dumps env vars, and emits CWD markers (see ``_wrap_command_script``)."""
|
||||
"""Full bash script: source snapshot, cd, run, re-dump env, emit CWD markers."""
|
||||
return _wrap_command_script(
|
||||
command,
|
||||
quoted_cwd=self._quote_cwd_for_cd(cwd),
|
||||
quoted_snap=self._quote_shell_path(self._snapshot_path),
|
||||
snap_tmp_template=self._quote_shell_path(self._snapshot_path + _SNAP_TMP_SUFFIX),
|
||||
passthrough_names=self._snapshot_excluded_passthrough_names(),
|
||||
snapshot_ready=self._snapshot_ready,
|
||||
cwd_marker=self._cwd_marker,
|
||||
)
|
||||
**self._snapshot_script_kwargs(cwd))
|
||||
|
||||
@staticmethod
|
||||
def _embed_stdin_heredoc(command: str, stdin_data: str) -> str:
|
||||
@@ -382,28 +315,19 @@ class BaseEnvironment(ABC):
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
def _wait_for_process(
|
||||
self,
|
||||
proc: ProcessHandle,
|
||||
timeout: int = 120,
|
||||
*,
|
||||
bounded_capture: bool = False,
|
||||
watch_interrupt_tid: int | None = None,
|
||||
) -> dict:
|
||||
"""Poll-based wait with interrupt checking and stdout draining.
|
||||
self, proc: ProcessHandle, timeout: int = 120, *,
|
||||
bounded_capture: bool = False, watch_interrupt_tid: int | None = None) -> dict:
|
||||
"""Poll-based wait with interrupt checking and stdout draining (shared, not overridden).
|
||||
|
||||
Shared across all backends — not overridden. ``bounded_capture=True``
|
||||
(foreground terminal-tool path only) retains at most
|
||||
``tool_output.max_bytes`` in a head/tail window so a verbose subprocess
|
||||
cannot OOM the process; the default keeps full fidelity for internal
|
||||
consumers (file-op ``cat`` reads, RPC reads) where truncation corrupts
|
||||
data. Fires the activity callback every 10s so the gateway's
|
||||
inactivity timeout doesn't kill long commands. ``watch_interrupt_tid``
|
||||
is the tool-worker thread that submitted this wait: ``execute()`` may
|
||||
move the wait onto a ``run_bounded_sync`` worker while ``/stop`` still
|
||||
interrupts the original tid, so both bits are honored. A
|
||||
``KeyboardInterrupt``/``SystemExit`` mid-poll kills the process first —
|
||||
the local backend spawns into its own process group (``os.setsid``), so
|
||||
an unkilled child would be orphaned to PPID=1 and keep running.
|
||||
``bounded_capture=True`` (foreground terminal-tool path only) retains at most
|
||||
``tool_output.max_bytes`` in a head/tail window so a verbose subprocess cannot OOM the
|
||||
process; the default keeps full fidelity for internal consumers where truncation
|
||||
corrupts data. Fires the activity callback every 10s so the gateway's inactivity
|
||||
timeout doesn't kill long commands. ``watch_interrupt_tid`` is the tool-worker thread
|
||||
that submitted this wait: ``execute()`` may move the wait onto a ``run_bounded_sync``
|
||||
worker while ``/stop`` still interrupts the original tid, so both bits are honored.
|
||||
``KeyboardInterrupt``/``SystemExit`` mid-poll kills the process first — the local
|
||||
backend spawns into its own process group, so an unkilled child would be orphaned.
|
||||
"""
|
||||
output = _new_output_collector(proc, bounded_capture)
|
||||
drain_thread = _start_drain_thread(proc, output)
|
||||
@@ -426,14 +350,11 @@ class BaseEnvironment(ABC):
|
||||
if is_interrupted() or is_thread_interrupted(watch_interrupt_tid):
|
||||
trace.interrupted()
|
||||
_kill_and_join()
|
||||
return self._finalize_wait_result(
|
||||
output, output.render(suffix="\n[Command interrupted]"), 130
|
||||
)
|
||||
return self._finalize_wait_result(output, output.render(suffix="\n[Command interrupted]"), 130)
|
||||
if time.monotonic() > deadline:
|
||||
trace.timed_out()
|
||||
_kill_and_join()
|
||||
timeout_msg = f"\n[Command timed out after {timeout}s]"
|
||||
rendered = output.render(suffix=timeout_msg)
|
||||
rendered = output.render(suffix=f"\n[Command timed out after {timeout}s]")
|
||||
if output.total_chars == 0:
|
||||
rendered = rendered.lstrip()
|
||||
return self._finalize_wait_result(output, rendered, 124)
|
||||
@@ -459,10 +380,9 @@ class BaseEnvironment(ABC):
|
||||
pass
|
||||
trace.natural_exit(proc.returncode)
|
||||
|
||||
# Join the stdin writer before reading its error list: a child that
|
||||
# exits without reading stdin can otherwise race ahead of a recorded
|
||||
# encode failure. The timeout is a pure safety net (write raises
|
||||
# BrokenPipeError once the pipe closes).
|
||||
# Join the stdin writer before reading its error list: a child that exits without
|
||||
# reading stdin can otherwise race ahead of a recorded encode failure. The timeout
|
||||
# is a pure safety net (write raises BrokenPipeError once the pipe closes).
|
||||
stdin_thread = getattr(proc, "_hermes_stdin_thread", None)
|
||||
if stdin_thread is not None:
|
||||
stdin_thread.join(timeout=5)
|
||||
@@ -493,16 +413,11 @@ class BaseEnvironment(ABC):
|
||||
self._extract_cwd_from_output(result)
|
||||
|
||||
def _extract_cwd_from_output(self, result: dict):
|
||||
"""Parse the ``__HERMES_CWD_{session}__`` marker from ``result["output"]``.
|
||||
|
||||
Updates ``self.cwd`` and strips the marker line. Sets
|
||||
``result["cwd_observed"]`` (and ``result["cwd"]``) only when THIS command
|
||||
emitted a marker: the wrapper prints it after the command returns, so a
|
||||
killed/timed-out command emits none and ``self.cwd`` keeps the previous
|
||||
value. The environment is shared across sessions, so callers must not
|
||||
attribute an unobserved cwd to this session, and concurrent callers
|
||||
must read ``result["cwd"]`` rather than ``self.cwd``.
|
||||
"""
|
||||
"""Parse the ``__HERMES_CWD_{session}__`` marker from ``result["output"]``, update
|
||||
``self.cwd`` and strip the marker line. ``result["cwd_observed"]``/``["cwd"]`` are set
|
||||
only when THIS command emitted a marker: a killed/timed-out command emits none and
|
||||
``self.cwd`` keeps the previous value. The environment is shared across sessions, so
|
||||
concurrent callers must read ``result["cwd"]`` rather than ``self.cwd``."""
|
||||
split = _split_cwd_marker(result.get("output", ""), self._cwd_marker)
|
||||
if split is None:
|
||||
return
|
||||
@@ -534,17 +449,15 @@ class BaseEnvironment(ABC):
|
||||
timeout: int | None = None,
|
||||
stdin_data: str | None = None,
|
||||
rewrite_compound_background: bool = True,
|
||||
bounded_capture: bool = False,
|
||||
) -> dict:
|
||||
bounded_capture: bool = False) -> dict:
|
||||
"""Execute a command, return {"output": str, "returncode": int}.
|
||||
|
||||
``bounded_capture=True`` caps retention at ``tool_output.max_bytes``
|
||||
WHILE draining; only the foreground terminal tool may set it — internal
|
||||
full-fidelity consumers (file-op ``cat`` reads feeding the patch
|
||||
engine, RPC reads, log reads) MUST leave it False or data is corrupted.
|
||||
The wait is bounded by ``agent.deadline.run_bounded_sync`` so a wedged
|
||||
poll loop cannot hang past ``timeout`` and silently disable every
|
||||
asyncio timer in the process.
|
||||
``bounded_capture=True`` caps retention at ``tool_output.max_bytes`` WHILE draining;
|
||||
only the foreground terminal tool may set it — internal full-fidelity consumers
|
||||
(file-op ``cat`` reads feeding the patch engine, RPC reads, log reads) MUST leave it
|
||||
False or data is corrupted. The wait is bounded by ``agent.deadline.run_bounded_sync``
|
||||
so a wedged poll loop cannot hang past ``timeout`` and silently disable every asyncio
|
||||
timer in the process.
|
||||
"""
|
||||
self._before_execute()
|
||||
|
||||
@@ -558,11 +471,7 @@ class BaseEnvironment(ABC):
|
||||
effective_cwd = cwd or self.cwd
|
||||
|
||||
# Merge sudo stdin with caller stdin.
|
||||
if sudo_stdin is not None:
|
||||
effective_stdin = sudo_stdin + (stdin_data or "")
|
||||
else:
|
||||
effective_stdin = stdin_data
|
||||
|
||||
effective_stdin = sudo_stdin + (stdin_data or "") if sudo_stdin is not None else stdin_data
|
||||
if effective_stdin and self._stdin_mode == "heredoc":
|
||||
exec_command = self._embed_stdin_heredoc(exec_command, effective_stdin)
|
||||
effective_stdin = None
|
||||
@@ -582,28 +491,20 @@ class BaseEnvironment(ABC):
|
||||
def _spawn_and_wait() -> dict:
|
||||
if parent_activity_cb is not None:
|
||||
set_activity_callback(parent_activity_cb)
|
||||
spawned = self._run_bash(
|
||||
wrapped, login=login, timeout=effective_timeout, stdin_data=effective_stdin
|
||||
)
|
||||
spawned = self._run_bash(wrapped, login=login, timeout=effective_timeout, stdin_data=effective_stdin)
|
||||
proc_holder.append(spawned)
|
||||
return self._wait_for_process(
|
||||
spawned,
|
||||
timeout=effective_timeout,
|
||||
bounded_capture=bounded_capture,
|
||||
watch_interrupt_tid=parent_tid,
|
||||
)
|
||||
spawned, timeout=effective_timeout, bounded_capture=bounded_capture, watch_interrupt_tid=parent_tid)
|
||||
|
||||
def _on_timeout() -> None:
|
||||
if not proc_holder:
|
||||
return
|
||||
self._kill_spawned_tree(proc_holder[0])
|
||||
if proc_holder:
|
||||
self._kill_spawned_tree(proc_holder[0])
|
||||
|
||||
# Hard wall-clock backstop: ``_wait_for_process`` polls to
|
||||
# ``effective_timeout`` on the tool thread; if that is the event-loop
|
||||
# thread, or the wait never returns (Windows pipe/poll hang), every
|
||||
# asyncio timer is silently disabled. ``run_bounded_sync`` drives expiry
|
||||
# from a daemon worker + ``Event.wait`` so a blocked loop cannot disable
|
||||
# it; the grace lets the inner loop return the partial-output 124 path.
|
||||
# Hard wall-clock backstop: ``_wait_for_process`` polls to ``effective_timeout`` on the
|
||||
# tool thread; if that is the event-loop thread, or the wait never returns (Windows
|
||||
# pipe/poll hang), every asyncio timer is silently disabled. ``run_bounded_sync`` drives
|
||||
# expiry from a daemon worker + ``Event.wait`` so a blocked loop cannot disable it; the
|
||||
# grace lets the inner loop return the partial-output 124 path.
|
||||
from agent.deadline import run_bounded_sync
|
||||
|
||||
try:
|
||||
@@ -614,11 +515,7 @@ class BaseEnvironment(ABC):
|
||||
|
||||
try:
|
||||
bounded = run_bounded_sync(
|
||||
_spawn_and_wait,
|
||||
bound_s,
|
||||
label=f"terminal.wait:{type(self).__name__}",
|
||||
on_timeout=_on_timeout,
|
||||
)
|
||||
_spawn_and_wait, bound_s, label=f"terminal.wait:{type(self).__name__}", on_timeout=_on_timeout)
|
||||
except (KeyboardInterrupt, SystemExit):
|
||||
_on_timeout()
|
||||
raise
|
||||
@@ -629,7 +526,6 @@ class BaseEnvironment(ABC):
|
||||
else:
|
||||
result = bounded.value
|
||||
self._update_cwd(result)
|
||||
|
||||
return result
|
||||
|
||||
def _kill_spawned_tree(self, spawned) -> None:
|
||||
@@ -643,7 +539,6 @@ class BaseEnvironment(ABC):
|
||||
return
|
||||
try:
|
||||
from agent.deadline import kill_process_tree
|
||||
|
||||
kill_process_tree(int(pid))
|
||||
except Exception:
|
||||
logger.debug("terminal wait-bound kill_process_tree failed", exc_info=True)
|
||||
@@ -665,5 +560,4 @@ class BaseEnvironment(ABC):
|
||||
def _prepare_command(self, command: str) -> tuple[str, str | None]:
|
||||
"""Transform sudo commands if SUDO_PASSWORD is available."""
|
||||
from tools.terminal_tool import _transform_sudo_command
|
||||
|
||||
return _transform_sudo_command(command)
|
||||
|
||||
@@ -26,12 +26,9 @@ _SPILL_MAX_AGE_S = 7 * 86400
|
||||
|
||||
|
||||
class _BoundedOutputCollector:
|
||||
"""Retain a bounded 40/60 head-tail window of streamed text.
|
||||
|
||||
When ``spill_path`` is set, the FULL stream is also teed to that file once
|
||||
eviction begins (up to ``_SPILL_CAP_CHARS``) so a truncated foreground
|
||||
result is recoverable without re-running. Memory stays bounded either way.
|
||||
"""
|
||||
"""Retain a bounded 40/60 head-tail window of streamed text. When ``spill_path`` is set,
|
||||
the FULL stream is also teed to that file once eviction begins (up to ``_SPILL_CAP_CHARS``)
|
||||
so a truncated foreground result is recoverable without re-running."""
|
||||
|
||||
# Hard ceiling on spill file size; protects disk from runaway output.
|
||||
_SPILL_CAP_CHARS = 5_000_000
|
||||
@@ -58,13 +55,10 @@ class _BoundedOutputCollector:
|
||||
try:
|
||||
if self._spill_fh is None:
|
||||
from tools.spill_safety import ensure_spill_dir, open_exclusive
|
||||
|
||||
# Raw pre-redaction output: private perms + symlink-refusing
|
||||
# exclusive create (a planted link must fail, never redirect).
|
||||
ensure_spill_dir(self._spill_path.parent, private=True)
|
||||
self._spill_fh = open_exclusive(
|
||||
self._spill_path, private=True, errors="replace"
|
||||
)
|
||||
self._spill_fh = open_exclusive(self._spill_path, private=True, errors="replace")
|
||||
# Backfill what's retained so the file holds the stream from byte 0.
|
||||
backlog = "".join(self._head) + "".join(self._tail)
|
||||
self._spill_fh.write(backlog)
|
||||
@@ -110,9 +104,7 @@ class _BoundedOutputCollector:
|
||||
text_len = len(text)
|
||||
# Spill tee activates at the first overflow, then mirrors every chunk.
|
||||
if self._spill_path is not None and (
|
||||
self._spill_fh is not None
|
||||
or self._total_chars + text_len > self.max_chars
|
||||
):
|
||||
self._spill_fh is not None or self._total_chars + text_len > self.max_chars):
|
||||
self._maybe_spill(text)
|
||||
self._total_chars += text_len
|
||||
start = 0
|
||||
@@ -164,12 +156,10 @@ class _BoundedOutputCollector:
|
||||
for _ in range(4):
|
||||
content_budget = max(0, available - len(notice))
|
||||
head_chars = int(content_budget * 0.4)
|
||||
tail_chars = content_budget - head_chars
|
||||
omitted = max(0, self._total_chars - head_chars - tail_chars)
|
||||
omitted = max(0, self._total_chars - content_budget)
|
||||
updated = (
|
||||
f"\n\n... [OUTPUT TRUNCATED - {omitted:,} chars omitted "
|
||||
f"out of {self._total_chars:,} total] ...\n\n"
|
||||
)
|
||||
f"out of {self._total_chars:,} total] ...\n\n")
|
||||
if updated == notice:
|
||||
break
|
||||
notice = updated
|
||||
@@ -182,20 +172,15 @@ class _BoundedOutputCollector:
|
||||
|
||||
|
||||
def _new_output_collector(proc, bounded_capture: bool) -> _BoundedOutputCollector:
|
||||
"""Build the collector for one ``_wait_for_process`` call.
|
||||
|
||||
``bounded_capture`` (foreground terminal path only) caps retention at
|
||||
``tool_output.max_bytes`` and tees overflow to a spill file under
|
||||
``$HERMES_HOME/cache/terminal-output`` (created only on actual overflow;
|
||||
spills older than 7 days are pruned opportunistically). Otherwise the
|
||||
collector is effectively unbounded so internal consumers (file-op ``cat``
|
||||
reads, RPC reads) keep full-fidelity output.
|
||||
"""
|
||||
"""Build the collector for one ``_wait_for_process`` call. ``bounded_capture`` (foreground
|
||||
terminal path only) caps retention at ``tool_output.max_bytes`` and tees overflow to a
|
||||
spill file under ``$HERMES_HOME/cache/terminal-output`` (created only on actual overflow;
|
||||
spills older than 7 days are pruned opportunistically). Otherwise the collector is
|
||||
effectively unbounded so internal consumers keep full-fidelity output."""
|
||||
if not bounded_capture:
|
||||
return _BoundedOutputCollector(_UNBOUNDED_CAPTURE_CHARS)
|
||||
try:
|
||||
from tools.tool_output_limits import get_max_bytes
|
||||
|
||||
capture_limit = get_max_bytes()
|
||||
except Exception:
|
||||
capture_limit = 50_000
|
||||
@@ -216,9 +201,7 @@ def _new_output_collector(proc, bounded_capture: bool) -> _BoundedOutputCollecto
|
||||
return _BoundedOutputCollector(capture_limit, spill_path=spill_path)
|
||||
|
||||
|
||||
def _finalize_wait_result(
|
||||
collector: _BoundedOutputCollector, rendered: str, returncode: int | None
|
||||
) -> dict:
|
||||
def _finalize_wait_result(collector: _BoundedOutputCollector, rendered: str, returncode: int | None) -> dict:
|
||||
"""Assemble a wait result, attaching spill metadata when overflow occurred."""
|
||||
result = {"output": rendered, "returncode": returncode}
|
||||
spill = collector.close_spill()
|
||||
@@ -236,16 +219,13 @@ def _finalize_wait_result(
|
||||
def _pipe_stdin(proc: subprocess.Popen, data: str) -> None:
|
||||
"""Write *data* to proc.stdin on a daemon thread to avoid pipe-buffer deadlocks.
|
||||
|
||||
Writes go through ``proc.stdin.buffer`` as UTF-8 bytes we encode ourselves:
|
||||
Windows text-mode stdin would translate ``\\n`` -> ``\\r\\n`` and corrupt
|
||||
every write_file/patch payload (byte-identical on POSIX). Encoding uses
|
||||
``surrogateescape`` (exact inverse of the read-side decode). Surrogates
|
||||
outside U+DC80-U+DCFF raise; the error is recorded on
|
||||
``proc._hermes_stdin_errors`` and stdin is still closed in ``finally`` so
|
||||
the child sees EOF instead of hanging. ``_wait_for_process`` surfaces the
|
||||
recorded error as ``stdin_error``.
|
||||
Writes go through ``proc.stdin.buffer`` as UTF-8 bytes we encode ourselves: Windows
|
||||
text-mode stdin would translate ``\\n`` -> ``\\r\\n`` and corrupt every write_file/patch
|
||||
payload. Encoding uses ``surrogateescape`` (exact inverse of the read-side decode);
|
||||
surrogates outside U+DC80-U+DCFF raise, the error is recorded on
|
||||
``proc._hermes_stdin_errors`` (surfaced by ``_wait_for_process`` as ``stdin_error``) and
|
||||
stdin is still closed in ``finally`` so the child sees EOF instead of hanging.
|
||||
"""
|
||||
|
||||
errors: list[BaseException] = []
|
||||
proc._hermes_stdin_errors = errors
|
||||
|
||||
@@ -278,15 +258,10 @@ def _pipe_stdin(proc: subprocess.Popen, data: str) -> None:
|
||||
thread.start()
|
||||
|
||||
|
||||
def _popen_bash(
|
||||
cmd: list[str], stdin_data: str | None = None, **kwargs
|
||||
) -> subprocess.Popen:
|
||||
"""Spawn a subprocess with standard stdout/stderr/stdin setup.
|
||||
|
||||
If *stdin_data* is provided, writes it asynchronously via :func:`_pipe_stdin`.
|
||||
Backends with special Popen needs (e.g. local's ``preexec_fn``) can bypass
|
||||
this and call :func:`_pipe_stdin` directly.
|
||||
"""
|
||||
def _popen_bash(cmd: list[str], stdin_data: str | None = None, **kwargs) -> subprocess.Popen:
|
||||
"""Spawn a subprocess with standard stdout/stderr/stdin setup; *stdin_data* is written
|
||||
asynchronously via :func:`_pipe_stdin`. Backends with special Popen needs (e.g. local's
|
||||
``preexec_fn``) can bypass this and call :func:`_pipe_stdin` directly."""
|
||||
kwargs.setdefault("creationflags", windows_hide_flags())
|
||||
proc = subprocess.Popen(
|
||||
cmd,
|
||||
@@ -294,8 +269,7 @@ def _popen_bash(
|
||||
stderr=subprocess.STDOUT,
|
||||
stdin=subprocess.PIPE if stdin_data is not None else subprocess.DEVNULL,
|
||||
text=True, encoding="utf-8", errors="replace",
|
||||
**kwargs,
|
||||
)
|
||||
**kwargs)
|
||||
if stdin_data is not None:
|
||||
_pipe_stdin(proc, stdin_data)
|
||||
return proc
|
||||
@@ -307,11 +281,8 @@ def _popen_bash(
|
||||
|
||||
|
||||
class ProcessHandle(Protocol):
|
||||
"""Duck type that every backend's _run_bash() must return.
|
||||
|
||||
subprocess.Popen satisfies this natively. SDK backends (Modal, Daytona)
|
||||
return _ThreadedProcessHandle which adapts their blocking calls.
|
||||
"""
|
||||
"""Duck type every backend's _run_bash() must return. subprocess.Popen satisfies this
|
||||
natively; SDK backends (Modal, Daytona) return _ThreadedProcessHandle."""
|
||||
|
||||
def poll(self) -> int | None: ...
|
||||
def kill(self) -> None: ...
|
||||
@@ -325,19 +296,11 @@ class ProcessHandle(Protocol):
|
||||
|
||||
|
||||
class _ThreadedProcessHandle:
|
||||
"""Adapter for SDK backends (Modal, Daytona) that have no real subprocess.
|
||||
"""Adapter for SDK backends (Modal, Daytona) that have no real subprocess: runs a blocking
|
||||
``exec_fn() -> (output_str, exit_code)`` on a background thread behind a ProcessHandle
|
||||
interface. ``cancel_fn`` is invoked on ``kill()`` for backend-specific cancellation."""
|
||||
|
||||
Wraps a blocking ``exec_fn() -> (output_str, exit_code)`` in a background
|
||||
thread and exposes a ProcessHandle-compatible interface. An optional
|
||||
``cancel_fn`` is invoked on ``kill()`` for backend-specific cancellation
|
||||
(e.g. Modal sandbox.terminate, Daytona sandbox.stop).
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
exec_fn: Callable[[], tuple[str, int]],
|
||||
cancel_fn: Callable[[], None] | None = None,
|
||||
):
|
||||
def __init__(self, exec_fn: Callable[[], tuple[str, int]], cancel_fn: Callable[[], None] | None = None):
|
||||
self._cancel_fn = cancel_fn
|
||||
self._done = threading.Event()
|
||||
self._returncode: int | None = None
|
||||
@@ -366,8 +329,7 @@ class _ThreadedProcessHandle:
|
||||
pass
|
||||
self._done.set()
|
||||
|
||||
t = threading.Thread(target=_worker, daemon=True)
|
||||
t.start()
|
||||
threading.Thread(target=_worker, daemon=True).start()
|
||||
|
||||
@property
|
||||
def stdout(self):
|
||||
@@ -400,20 +362,15 @@ class _ThreadedProcessHandle:
|
||||
def _drain_stdout(proc: ProcessHandle, output: _BoundedOutputCollector) -> None:
|
||||
"""Drain ``proc.stdout`` into *output* until EOF or shortly after exit.
|
||||
|
||||
``for line in proc.stdout`` would block on ``readline()`` until EOF, and a
|
||||
backgrounded grandchild (``cmd &``, ``setsid cmd & disown``) inherits the
|
||||
pipe's write end — so the tool would hang for the grandchild's lifetime.
|
||||
Instead we ``select()`` with a short poll and stop ~300ms after bash exits
|
||||
even if the pipe has not EOF'd (later grandchild output goes to an orphaned
|
||||
pipe, harmless). Raw 4096-byte ``os.read`` chunks can split a multibyte
|
||||
UTF-8 sequence, so an incremental decoder with ``errors="replace"`` (same
|
||||
as the Popen TextIOWrapper) buffers partial sequences across chunks.
|
||||
|
||||
Streams without a real integer ``fileno()`` (mocks, iterator-style stdout
|
||||
from in-memory adapters) are iterated to EOF instead — otherwise the thread
|
||||
would die silently and lose all output. ``select()`` does not work on pipe
|
||||
fds on Windows, so a blocking ``os.read`` loop is used there (EOF arrives
|
||||
promptly when bash exits).
|
||||
``for line in proc.stdout`` would block on ``readline()`` until EOF, and a backgrounded
|
||||
grandchild (``cmd &``, ``setsid cmd & disown``) inherits the pipe's write end — so the
|
||||
tool would hang for the grandchild's lifetime. Instead we ``select()`` with a short poll
|
||||
and stop ~300ms after bash exits even if the pipe has not EOF'd. Raw 4096-byte ``os.read``
|
||||
chunks can split a multibyte UTF-8 sequence, so an incremental decoder with
|
||||
``errors="replace"`` buffers partial sequences across chunks. Streams without a real
|
||||
integer ``fileno()`` (mocks, in-memory adapters) are iterated to EOF instead — otherwise
|
||||
the thread would die silently and lose all output. ``select()`` does not work on pipe fds
|
||||
on Windows, so a blocking ``os.read`` loop is used there.
|
||||
"""
|
||||
stream = proc.stdout
|
||||
if stream is None:
|
||||
|
||||
@@ -9,17 +9,15 @@ import re
|
||||
import shlex
|
||||
from typing import Iterable
|
||||
|
||||
# Bridged per-session vars (gateway.session_context._VAR_MAP) are injected fresh
|
||||
# onto every command's process env and must NEVER persist in the shared bash
|
||||
# snapshot: one long-lived backend serves many sessions, so a snapshot carrying
|
||||
# the FIRST session's HERMES_SESSION_ID would make every LATER session source a
|
||||
# foreign identity, overriding the correct per-command Popen env. Every bridged
|
||||
# name starts with one of these prefixes (or is HERMES_UI_SESSION_ID); unit
|
||||
# tests use this regex as the Python-side contract for the exclusion set.
|
||||
# Bridged per-session vars (gateway.session_context._VAR_MAP) are injected fresh onto every
|
||||
# command's process env and must NEVER persist in the shared bash snapshot: one long-lived
|
||||
# backend serves many sessions, so a snapshot carrying the FIRST session's HERMES_SESSION_ID
|
||||
# would make every LATER session source a foreign identity. Every bridged name starts with
|
||||
# one of these prefixes (or is HERMES_UI_SESSION_ID); unit tests use this regex as the
|
||||
# Python-side contract for the exclusion set.
|
||||
_SNAPSHOT_EXCLUDED_ENV_REGEX = (
|
||||
"^declare -x (HERMES_SESSION_|HERMES_UI_SESSION_ID|HERMES_CRON_AUTO_DELIVER_|"
|
||||
"HERMES_CRON_SESSION|HERMES_BROWSER_CONTROL_)"
|
||||
)
|
||||
"HERMES_CRON_SESSION|HERMES_BROWSER_CONTROL_)")
|
||||
_SHELL_ENV_NAME_RE = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
|
||||
|
||||
# mktemp template suffix + the shell variable holding the allocated temp path.
|
||||
@@ -37,31 +35,23 @@ def _cwd_marker_printf(marker: str) -> str:
|
||||
return f"printf '\\n{marker}%s{marker}\\n' \"$(pwd -P)\""
|
||||
|
||||
|
||||
def _export_dump_excluding_session_vars(
|
||||
tmp_path: str,
|
||||
excluded_names: Iterable[str] = (),
|
||||
) -> str:
|
||||
"""Shell snippet dumping ``export -p`` to *tmp_path* minus the per-session
|
||||
bridged vars (see ``_SNAPSHOT_EXCLUDED_ENV_REGEX``) and *excluded_names*.
|
||||
def _export_dump_excluding_session_vars(tmp_path: str, excluded_names: Iterable[str] = ()) -> str:
|
||||
"""Shell snippet dumping ``export -p`` to *tmp_path* minus the per-session bridged vars
|
||||
(see ``_SNAPSHOT_EXCLUDED_ENV_REGEX``) and *excluded_names*.
|
||||
|
||||
The vars are ``unset`` in a subshell BEFORE ``export -p``. A line-based
|
||||
``grep -vE`` filter is unsafe: bash 3.2 prints a value containing a newline
|
||||
as a multi-line ``declare -x`` block, so continuation lines (attacker text
|
||||
smuggled via e.g. HERMES_SESSION_CHAT_NAME) would survive into the snapshot
|
||||
and execute on the next ``source``. ``|| true`` keeps the success contract.
|
||||
|
||||
The dump is a brace group with the redirection on the group: *tmp_path* is
|
||||
usually a shell-variable expansion, and a redirect attached to a pipeline
|
||||
segment would expand it inside that segment's subshell, inconsistently with
|
||||
The vars are ``unset`` in a subshell BEFORE ``export -p``: a line-based ``grep -vE`` is
|
||||
unsafe because bash 3.2 prints a value containing a newline as a multi-line ``declare -x``
|
||||
block, so smuggled continuation lines would survive into the snapshot and execute on the
|
||||
next ``source``. ``|| true`` keeps the success contract. The dump is a brace group with the
|
||||
redirection on the group: *tmp_path* is usually a shell-variable expansion, and a redirect
|
||||
on a pipeline segment would expand it inside that segment's subshell, inconsistently with
|
||||
the parent that expands the follow-up ``mv``.
|
||||
"""
|
||||
# ${!PREFIX*} is bash 3.2+ name-prefix expansion; empty matches are ignored
|
||||
# under 2>/dev/null. Caller names are quoted so malformed config can never
|
||||
# become shell syntax (valid names stay unquoted by shlex.quote()).
|
||||
safe_names = {name for name in excluded_names if isinstance(name, str) and name}
|
||||
extra_unset = " ".join(shlex.quote(name) for name in sorted(safe_names))
|
||||
if extra_unset:
|
||||
extra_unset = f" {extra_unset}"
|
||||
extra_unset = "".join(f" {shlex.quote(name)}" for name in sorted(safe_names))
|
||||
return (
|
||||
"{ ( "
|
||||
"unset ${!HERMES_SESSION_*} ${!HERMES_CRON_AUTO_DELIVER_*} "
|
||||
@@ -73,30 +63,22 @@ def _export_dump_excluding_session_vars(
|
||||
f"HERMES_UI_SESSION_ID{extra_unset} 2>/dev/null; "
|
||||
"export -p; "
|
||||
") || true; } "
|
||||
f"> {tmp_path}"
|
||||
)
|
||||
f"> {tmp_path}")
|
||||
|
||||
|
||||
def _snapshot_bootstrap_script(
|
||||
*,
|
||||
quoted_cwd: str,
|
||||
quoted_snap: str,
|
||||
snap_tmp_template: str,
|
||||
excluded_names: Iterable[str],
|
||||
cwd_marker: str,
|
||||
*, quoted_cwd: str, quoted_snap: str, snap_tmp_template: str, excluded_names: Iterable[str], cwd_marker: str,
|
||||
) -> str:
|
||||
"""Login-shell bootstrap that captures env/functions/aliases into the snapshot.
|
||||
|
||||
Atomic publish: assemble in a ``mktemp`` file, then ``mv`` over the final
|
||||
path so a concurrent ``source`` never reads a half-written snapshot. The
|
||||
temp name must be unique per concurrent writer: ``$$`` is the parent PID in
|
||||
``&``-launched subshells and macOS bash 3.2 lacks ``$BASHPID``, so only
|
||||
``mktemp`` is portable. Functions are filtered by NAME via ``declare -F``
|
||||
(a line-based ``declare -f | grep -v`` strips the header and leaves an
|
||||
orphaned body that breaks every sourced command); the non-empty guard
|
||||
matters because bare ``declare -f`` dumps ALL functions. The trailing
|
||||
``cd`` restores the configured cwd after profile scripts (e.g. ``cd ~``)
|
||||
so ``pwd -P`` reports terminal.cwd, not the profile's directory.
|
||||
Atomic publish: assemble in a ``mktemp`` file, then ``mv`` over the final path so a
|
||||
concurrent ``source`` never reads a half-written snapshot (``$$`` is the parent PID in
|
||||
``&``-launched subshells and macOS bash 3.2 lacks ``$BASHPID``, so only ``mktemp`` is
|
||||
portable). Functions are filtered by NAME via ``declare -F`` (a line-based
|
||||
``declare -f | grep -v`` strips the header and leaves an orphaned body that breaks every
|
||||
sourced command); the non-empty guard matters because bare ``declare -f`` dumps ALL
|
||||
functions. The trailing ``cd`` restores the configured cwd after profile scripts
|
||||
(e.g. ``cd ~``) so ``pwd -P`` reports terminal.cwd, not the profile's directory.
|
||||
"""
|
||||
return (
|
||||
"umask 077\n"
|
||||
@@ -111,91 +93,68 @@ def _snapshot_bootstrap_script(
|
||||
# Publish only if assembly succeeded; otherwise drop the partial temp.
|
||||
f"mv -f {_SNAP_TMP} {quoted_snap} || rm -f {_SNAP_TMP}\n"
|
||||
f"builtin cd -- {quoted_cwd} 2>/dev/null || true\n"
|
||||
f"{_cwd_marker_printf(cwd_marker)}\n"
|
||||
)
|
||||
f"{_cwd_marker_printf(cwd_marker)}\n")
|
||||
|
||||
|
||||
def _passthrough_save_restore(names: Iterable[str]) -> tuple[list[str], list[str]]:
|
||||
"""Shell lines that save profile-scoped passthrough vars before the snapshot
|
||||
is sourced and restore (or unset) them afterwards.
|
||||
|
||||
A shared snapshot may hold the previous profile's value. Values stay in
|
||||
environment memory and never enter the command string, so secrets are not
|
||||
exposed through process arguments or logs.
|
||||
"""
|
||||
"""Shell lines that save profile-scoped passthrough vars before the snapshot is sourced
|
||||
and restore (or unset) them afterwards — a shared snapshot may hold the previous
|
||||
profile's value. Values stay in environment memory and never enter the command string."""
|
||||
save: list[str] = []
|
||||
restore: list[str] = []
|
||||
for name in names:
|
||||
marker = f"_HERMES_RUNTIME_PASSTHROUGH_{name}"
|
||||
present, value = f"{marker}_PRESENT", f"{marker}_VALUE"
|
||||
save.append(f"{present}=${{{name}+x}}")
|
||||
save.append(f"{value}=${{{name}-}}")
|
||||
restore.append(
|
||||
f'if [ "${present}" = x ]; then export {name}="${value}"; '
|
||||
f'else unset {name}; fi'
|
||||
)
|
||||
restore.append(f"unset {present} {value}")
|
||||
save += [f"{present}=${{{name}+x}}", f"{value}=${{{name}-}}"]
|
||||
restore += [
|
||||
f'if [ "${present}" = x ]; then export {name}="${value}"; else unset {name}; fi',
|
||||
f"unset {present} {value}"]
|
||||
return save, restore
|
||||
|
||||
|
||||
def _wrap_command_script(
|
||||
command: str,
|
||||
*,
|
||||
quoted_cwd: str,
|
||||
quoted_snap: str,
|
||||
snap_tmp_template: str,
|
||||
passthrough_names: Iterable[str],
|
||||
snapshot_ready: bool,
|
||||
cwd_marker: str,
|
||||
) -> str:
|
||||
command: str, *, quoted_cwd: str, quoted_snap: str, snap_tmp_template: str,
|
||||
passthrough_names: Iterable[str], snapshot_ready: bool, cwd_marker: str) -> str:
|
||||
"""Per-command bash script: source snapshot, cd, run, re-dump env, emit CWD marker.
|
||||
|
||||
``source`` stdout goes to /dev/null because macOS bash 3.2 / some Homebrew
|
||||
builds echo ``declare -x`` lines when sourcing. AI_AGENT/HERMES_AGENT
|
||||
advertise the harness to remote backends (whose env is not inherited from
|
||||
the Hermes process); ``${VAR:-default}`` never clobbers an outer harness.
|
||||
GIT_PAGER/PAGER=cat stop pager-happy tools hanging a PTY-backed command.
|
||||
The env re-dump uses the same mktemp+mv atomic publish as the bootstrap and
|
||||
chains ``mv`` on the dump succeeding so a failed dump never replaces a good
|
||||
snapshot. ``umask 077`` is applied after the user's command so snapshot
|
||||
files (which may carry secrets) are private without changing the command's
|
||||
umask.
|
||||
``source`` stdout goes to /dev/null because macOS bash 3.2 / some Homebrew builds echo
|
||||
``declare -x`` lines when sourcing. AI_AGENT/HERMES_AGENT advertise the harness to remote
|
||||
backends (whose env is not inherited); ``${VAR:-default}`` never clobbers an outer harness.
|
||||
GIT_PAGER/PAGER=cat stop pager-happy tools hanging a PTY-backed command. The env re-dump
|
||||
uses the same mktemp+mv atomic publish as the bootstrap and chains ``mv`` on the dump
|
||||
succeeding so a failed dump never replaces a good snapshot. ``umask 077`` is applied after
|
||||
the user's command so snapshot files (which may carry secrets) are private without
|
||||
changing the command's umask.
|
||||
"""
|
||||
escaped = command.replace("'", "'\\''")
|
||||
save, restore = _passthrough_save_restore(passthrough_names)
|
||||
parts = list(save)
|
||||
if snapshot_ready:
|
||||
parts.append(f"source {quoted_snap} >/dev/null 2>&1 || true")
|
||||
parts.extend(restore)
|
||||
parts.append(
|
||||
'export AI_AGENT="${AI_AGENT:-hermes-agent}" '
|
||||
'HERMES_AGENT="${HERMES_AGENT:-true}"'
|
||||
)
|
||||
parts.append('export GIT_PAGER="${GIT_PAGER:-cat}" PAGER="${PAGER:-cat}"')
|
||||
# ``--`` keeps hyphen-prefixed directory names from being parsed as options.
|
||||
parts.append(f"builtin cd -- {quoted_cwd} || exit 126")
|
||||
parts.append(f"eval '{escaped}'")
|
||||
parts.append("__hermes_ec=$?")
|
||||
parts.append("umask 077")
|
||||
parts += restore
|
||||
parts += [
|
||||
'export AI_AGENT="${AI_AGENT:-hermes-agent}" HERMES_AGENT="${HERMES_AGENT:-true}"',
|
||||
'export GIT_PAGER="${GIT_PAGER:-cat}" PAGER="${PAGER:-cat}"',
|
||||
# ``--`` keeps hyphen-prefixed directory names from being parsed as options.
|
||||
f"builtin cd -- {quoted_cwd} || exit 126",
|
||||
f"eval '{escaped}'",
|
||||
"__hermes_ec=$?",
|
||||
"umask 077"]
|
||||
if snapshot_ready:
|
||||
parts.append(
|
||||
f"__hermes_snap_tmp=$(mktemp {snap_tmp_template}) && "
|
||||
f"{{ {_export_dump_excluding_session_vars(_SNAP_TMP, passthrough_names)} "
|
||||
f"&& mv -f {_SNAP_TMP} {quoted_snap}; }} "
|
||||
f"2>/dev/null || rm -f {_SNAP_TMP} 2>/dev/null || true"
|
||||
)
|
||||
parts.append(_cwd_marker_printf(cwd_marker))
|
||||
parts.append("exit $__hermes_ec")
|
||||
f"2>/dev/null || rm -f {_SNAP_TMP} 2>/dev/null || true")
|
||||
parts += [_cwd_marker_printf(cwd_marker), "exit $__hermes_ec"]
|
||||
return "\n".join(parts)
|
||||
|
||||
|
||||
def _split_cwd_marker(output: str, marker: str) -> tuple[str | None, str] | None:
|
||||
"""Locate the last ``marker<path>marker`` pair in *output*.
|
||||
|
||||
Returns ``(cwd_path_or_None, output_without_marker_line)``, or ``None`` when
|
||||
no complete pair exists. The stripped span runs from the ``\\n`` the wrapper
|
||||
injected before the marker through the end of the marker line.
|
||||
"""
|
||||
"""Locate the last ``marker<path>marker`` pair in *output*. Returns
|
||||
``(cwd_path_or_None, output_without_marker_line)``, or ``None`` when no complete pair
|
||||
exists. The stripped span runs from the ``\\n`` the wrapper injected before the marker
|
||||
through the end of the marker line."""
|
||||
last = output.rfind(marker)
|
||||
if last == -1:
|
||||
return None
|
||||
|
||||
@@ -1,9 +1,6 @@
|
||||
"""Opt-in interrupt/poll tracing for ``BaseEnvironment._wait_for_process``.
|
||||
|
||||
Enabled by ``HERMES_DEBUG_INTERRUPT=1`` (see ``tools.environments.base``):
|
||||
logs loop entry/exit, interrupt/timeout detection, and ~30s heartbeats so
|
||||
"agent never sees the interrupt" reports can be diagnosed from agent.log.
|
||||
"""
|
||||
"""Opt-in interrupt/poll tracing for ``BaseEnvironment._wait_for_process``
|
||||
(``HERMES_DEBUG_INTERRUPT=1``): loop entry/exit, interrupt/timeout detection and ~30s
|
||||
heartbeats so "agent never sees the interrupt" reports can be diagnosed from agent.log."""
|
||||
|
||||
import threading
|
||||
import time
|
||||
@@ -16,8 +13,7 @@ _FMT = {
|
||||
"TIMEOUT": "iter=%d timeout=%ss",
|
||||
"HEARTBEAT": "iter=%d elapsed=%.0fs interrupt=%s activity_cb=%s%s",
|
||||
"EXCEPTION_EXIT": "iter=%d elapsed=%.1fs — killing subprocess group before re-raise",
|
||||
"EXIT (natural)": "iter=%d elapsed=%.1fs returncode=%s",
|
||||
}
|
||||
"EXIT (natural)": "iter=%d elapsed=%.1fs returncode=%s"}
|
||||
|
||||
|
||||
class _WaitTrace:
|
||||
@@ -39,8 +35,7 @@ class _WaitTrace:
|
||||
def _emit(self, event: str, *args) -> None:
|
||||
self._log.info(
|
||||
"[interrupt-debug] _wait_for_process %s tid=%s pid=%s " + _FMT[event],
|
||||
event, self._tid, self._pid, *args,
|
||||
)
|
||||
event, self._tid, self._pid, *args)
|
||||
|
||||
def _elapsed(self) -> float:
|
||||
return time.monotonic() - self._start
|
||||
@@ -67,8 +62,7 @@ class _WaitTrace:
|
||||
self._emit(
|
||||
"HEARTBEAT", self.iterations, self._elapsed(), is_interrupted(),
|
||||
"set" if not cb_now_none else "MISSING",
|
||||
" (LOST during run)" if cb_now_none and not self._cb_was_none else "",
|
||||
)
|
||||
" (LOST during run)" if cb_now_none and not self._cb_was_none else "")
|
||||
self._last_heartbeat = time.monotonic()
|
||||
self._cb_was_none = cb_now_none
|
||||
|
||||
|
||||
+121
-203
@@ -20,11 +20,7 @@ import uuid
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
|
||||
from tools.environments.base import (
|
||||
BaseEnvironment,
|
||||
EnvironmentConnectionError,
|
||||
_popen_bash,
|
||||
)
|
||||
from tools.environments.base import BaseEnvironment, EnvironmentConnectionError, _popen_bash
|
||||
from tools.environments.docker_egress import ( # noqa: F401 — re-exported for tests/patch targets
|
||||
_EGRESS_LABEL_KEY,
|
||||
_critical_egress_env_names,
|
||||
@@ -35,25 +31,19 @@ from tools.environments.docker_egress import ( # noqa: F401 — re-exported for
|
||||
check_docker_env_collisions,
|
||||
check_extra_args_collisions,
|
||||
check_forward_env_collisions,
|
||||
merge_egress_env,
|
||||
)
|
||||
merge_egress_env)
|
||||
from tools.environments.path_utils import sanitize_task_id_for_path
|
||||
from tools.environments.remote_common import bash_argv, run_capture
|
||||
from tools.environments.local import (
|
||||
_HERMES_PROVIDER_ENV_BLOCKLIST,
|
||||
_is_hermes_internal_secret,
|
||||
)
|
||||
from tools.environments.local import _HERMES_PROVIDER_ENV_BLOCKLIST, _is_hermes_internal_secret
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# Docker Desktop install paths checked when 'docker' is not in PATH
|
||||
# (macOS Intel / Apple Silicon Homebrew / app bundle).
|
||||
_DOCKER_SEARCH_PATHS = [
|
||||
"/usr/local/bin/docker",
|
||||
"/opt/homebrew/bin/docker",
|
||||
"/Applications/Docker.app/Contents/Resources/bin/docker",
|
||||
]
|
||||
"/Applications/Docker.app/Contents/Resources/bin/docker"]
|
||||
|
||||
_docker_executable: Optional[str] = None # resolved once, cached
|
||||
_ENV_VAR_NAME_RE = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
|
||||
@@ -62,7 +52,6 @@ _ENV_VAR_NAME_RE = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
|
||||
def _normalize_forward_env_names(forward_env: list[str] | None) -> list[str]:
|
||||
"""Return a deduplicated list of valid environment variable names."""
|
||||
normalized: list[str] = []
|
||||
seen: set[str] = set()
|
||||
for item in forward_env or []:
|
||||
if not isinstance(item, str):
|
||||
logger.warning("Ignoring non-string docker_forward_env entry: %r", item)
|
||||
@@ -72,9 +61,7 @@ def _normalize_forward_env_names(forward_env: list[str] | None) -> list[str]:
|
||||
continue
|
||||
if not _ENV_VAR_NAME_RE.match(key):
|
||||
logger.warning("Ignoring invalid docker_forward_env entry: %r", item)
|
||||
continue
|
||||
if key not in seen:
|
||||
seen.add(key)
|
||||
elif key not in normalized:
|
||||
normalized.append(key)
|
||||
return normalized
|
||||
|
||||
@@ -91,13 +78,10 @@ def _normalize_env_dict(env: dict | None) -> dict[str, str]:
|
||||
if not isinstance(key, str) or not _ENV_VAR_NAME_RE.match(key.strip()):
|
||||
logger.warning("Ignoring invalid docker_env key: %r", key)
|
||||
continue
|
||||
key = key.strip()
|
||||
if not isinstance(value, str):
|
||||
if not isinstance(value, (int, float, bool)):
|
||||
logger.warning("Ignoring non-string docker_env value for %r: %r", key, value)
|
||||
continue
|
||||
value = str(value)
|
||||
normalized[key] = value
|
||||
if not isinstance(value, (str, int, float, bool)):
|
||||
logger.warning("Ignoring non-string docker_env value for %r: %r", key.strip(), value)
|
||||
continue
|
||||
normalized[key.strip()] = value if isinstance(value, str) else str(value)
|
||||
return normalized
|
||||
|
||||
|
||||
@@ -116,10 +100,7 @@ _LABEL_VALUE_OK_RE = re.compile(r"[^A-Za-z0-9_.-]")
|
||||
|
||||
|
||||
def _sanitize_label_value(value: str) -> str:
|
||||
"""Coerce *value* into a Docker label-safe form; empty/invalid input becomes ``"unknown"``.
|
||||
|
||||
Lossy — never round-trip a sanitized value back into application logic.
|
||||
"""
|
||||
"""Lossy coercion into a Docker label-safe form; empty/invalid input becomes ``"unknown"``."""
|
||||
if not isinstance(value, str) or not value:
|
||||
return "unknown"
|
||||
return _LABEL_VALUE_OK_RE.sub("_", value)[:63] or "unknown"
|
||||
@@ -131,11 +112,8 @@ _sandbox_dir_name = sanitize_task_id_for_path
|
||||
|
||||
|
||||
def _get_active_profile_name() -> str:
|
||||
"""Active Hermes profile name, or ``"default"`` on any error.
|
||||
|
||||
Resolved at container-create time so a container stays tagged with the
|
||||
profile that created it even if the process switches profiles later.
|
||||
"""
|
||||
"""Active Hermes profile name, or ``"default"`` on any error. Resolved at container-create
|
||||
time so a container stays tagged with its creator even if the process switches profiles."""
|
||||
try:
|
||||
from hermes_cli.profiles import get_active_profile_name
|
||||
return get_active_profile_name() or "default"
|
||||
@@ -146,11 +124,10 @@ def _get_active_profile_name() -> str:
|
||||
def _container_identity(shared_key: str = "") -> str:
|
||||
"""Profile label used for container reuse and orphan reaping.
|
||||
|
||||
Profiles are isolated by default; an explicit shared key lets trusted
|
||||
profiles share one Docker identity. Label sanitization is lossy and reuse
|
||||
is label-keyed, so two DIFFERENT keys colliding after sanitization would
|
||||
attach to the same container — a digest of the raw key disambiguates.
|
||||
Plain profile names keep their historical un-suffixed labels.
|
||||
Profiles are isolated by default; an explicit shared key lets trusted profiles share
|
||||
one Docker identity. Label sanitization is lossy and reuse is label-keyed, so a digest
|
||||
of the raw key disambiguates colliding keys. Plain profile names keep their historical
|
||||
un-suffixed labels.
|
||||
"""
|
||||
if not shared_key:
|
||||
return _sanitize_label_value(_get_active_profile_name())
|
||||
@@ -208,15 +185,12 @@ def _container_finished_at(docker_exe: str, container_id: str):
|
||||
missing/unparseable or Docker's never-finished zero value ``0001-01-01T00:00:00Z``."""
|
||||
try:
|
||||
result = run_capture(
|
||||
[docker_exe, "inspect", "--format", "{{.State.FinishedAt}}", container_id], timeout=10,
|
||||
)
|
||||
[docker_exe, "inspect", "--format", "{{.State.FinishedAt}}", container_id], timeout=10)
|
||||
except (subprocess.TimeoutExpired, OSError) as e:
|
||||
logger.debug("orphan reaper docker inspect %s failed: %s", container_id[:12], e)
|
||||
return None
|
||||
if result.returncode != 0:
|
||||
return None
|
||||
raw = result.stdout.strip()
|
||||
if not raw or raw.startswith("0001-01-01"):
|
||||
if result.returncode != 0 or not raw or raw.startswith("0001-01-01"):
|
||||
return None
|
||||
# Docker emits RFC3339 with nanoseconds; fromisoformat only takes microseconds.
|
||||
raw = re.sub(r"(\.\d{6})\d+", r"\1", raw).replace("Z", "+00:00")
|
||||
@@ -241,23 +215,16 @@ def find_docker() -> Optional[str]:
|
||||
override = os.getenv("HERMES_DOCKER_BINARY")
|
||||
if override and _is_executable(override):
|
||||
logger.info("Using HERMES_DOCKER_BINARY override: %s", override)
|
||||
_docker_executable = override
|
||||
return override
|
||||
found = shutil.which("docker")
|
||||
if found:
|
||||
_docker_executable = found
|
||||
return found
|
||||
found = shutil.which("podman")
|
||||
if found:
|
||||
found = override
|
||||
elif found := shutil.which("docker"):
|
||||
pass
|
||||
elif found := shutil.which("podman"):
|
||||
logger.info("Using podman as container runtime: %s", found)
|
||||
elif found := next((p for p in _DOCKER_SEARCH_PATHS if _is_executable(p)), None):
|
||||
logger.info("Found docker at non-PATH location: %s", found)
|
||||
if found:
|
||||
_docker_executable = found
|
||||
return found
|
||||
for path in _DOCKER_SEARCH_PATHS:
|
||||
if _is_executable(path):
|
||||
logger.info("Found docker at non-PATH location: %s", path)
|
||||
_docker_executable = path
|
||||
return path
|
||||
return None
|
||||
return found
|
||||
|
||||
|
||||
# Security flags applied to every container. The container is the security
|
||||
@@ -274,8 +241,7 @@ _BASE_SECURITY_ARGS = [
|
||||
"--cap-add", "FOWNER",
|
||||
"--security-opt", "no-new-privileges",
|
||||
"--tmpfs", "/tmp:rw,nosuid,size=512m",
|
||||
"--tmpfs", "/var/tmp:rw,noexec,nosuid,size=256m",
|
||||
]
|
||||
"--tmpfs", "/var/tmp:rw,noexec,nosuid,size=256m"]
|
||||
|
||||
_DEFAULT_PIDS_LIMIT = "256" # applied only when the pids cgroup controller is available
|
||||
|
||||
@@ -289,8 +255,7 @@ def _extra_args_set_shm_size(extra_args: list) -> bool:
|
||||
"""True when docker_extra_args already set ``--shm-size`` (then our default is skipped)."""
|
||||
return any(
|
||||
isinstance(a, str) and (a == "--shm-size" or a.startswith("--shm-size="))
|
||||
for a in (extra_args or [])
|
||||
)
|
||||
for a in (extra_args or []))
|
||||
|
||||
|
||||
# /run is separate from _BASE_SECURITY_ARGS: s6-overlay images exec
|
||||
@@ -309,9 +274,7 @@ _S6_INIT_ENTRYPOINTS = ("/init", "/package/admin/s6-overlay/command/init")
|
||||
def _build_security_args(run_as_host_user: bool, run_exec: bool = False) -> list[str]:
|
||||
"""Security/cap/tmpfs args for the privilege mode; ``run_exec`` mounts /run exec for s6 images."""
|
||||
args = list(_BASE_SECURITY_ARGS) + list(_RUN_TMPFS_EXEC if run_exec else _RUN_TMPFS_NOEXEC)
|
||||
if run_as_host_user:
|
||||
return args
|
||||
return args + list(_PRIVDROP_CAP_ARGS)
|
||||
return args if run_as_host_user else args + list(_PRIVDROP_CAP_ARGS)
|
||||
|
||||
|
||||
def _image_uses_init_entrypoint(docker_exe: str, image: str) -> bool:
|
||||
@@ -320,17 +283,14 @@ def _image_uses_init_entrypoint(docker_exe: str, image: str) -> bool:
|
||||
pulled, returns False and keeps the hardened defaults."""
|
||||
try:
|
||||
result = run_capture(
|
||||
[docker_exe, "image", "inspect", image, "--format", "{{json .Config.Entrypoint}}"],
|
||||
timeout=15,
|
||||
)
|
||||
[docker_exe, "image", "inspect", image, "--format", "{{json .Config.Entrypoint}}"], timeout=15)
|
||||
except (subprocess.SubprocessError, OSError) as e:
|
||||
logger.debug("Docker: could not inspect entrypoint for %s: %s", image, e)
|
||||
return False
|
||||
if result.returncode != 0:
|
||||
logger.debug(
|
||||
"Docker: image inspect for %s returned %d (stderr=%s)",
|
||||
image, result.returncode, result.stderr.strip(),
|
||||
)
|
||||
image, result.returncode, result.stderr.strip())
|
||||
return False
|
||||
raw = (result.stdout or "").strip()
|
||||
if not raw or raw == "null":
|
||||
@@ -381,8 +341,7 @@ def _cgroup_limits_available(image: str) -> bool:
|
||||
result = run_capture(
|
||||
[docker_exe, "run", "--rm", "--cpus", "0.5", "--memory", "64m", "--pids-limit", "32",
|
||||
image, "sleep", "0"],
|
||||
timeout=60,
|
||||
)
|
||||
timeout=60)
|
||||
_cgroup_limits_ok = result.returncode == 0
|
||||
if not _cgroup_limits_ok:
|
||||
logger.warning(
|
||||
@@ -391,8 +350,7 @@ def _cgroup_limits_available(image: str) -> bool:
|
||||
"CPU, memory or PID limits. To enable, delegate the cpu, "
|
||||
"memory and pids cgroup controllers to this container. "
|
||||
"Probe stderr: %s",
|
||||
(result.stderr or "").strip()[:500],
|
||||
)
|
||||
(result.stderr or "").strip()[:500])
|
||||
except Exception as e:
|
||||
_cgroup_limits_ok = False
|
||||
logger.warning("Cgroup limit probe failed; disabling resource limits: %s", e)
|
||||
@@ -414,8 +372,7 @@ def _ensure_docker_available() -> None:
|
||||
"CLI is available.",
|
||||
error="Docker executable not found in PATH or known install locations. "
|
||||
"Install Docker and ensure the 'docker' command is available.",
|
||||
hint="Install Docker (or fix PATH) and retry, or switch terminal.backend to 'local'.",
|
||||
)
|
||||
hint="Install Docker (or fix PATH) and retry, or switch terminal.backend to 'local'.")
|
||||
try:
|
||||
result = run_capture([docker_exe, "version"], timeout=5)
|
||||
except FileNotFoundError:
|
||||
@@ -423,16 +380,14 @@ def _ensure_docker_available() -> None:
|
||||
"Docker backend selected but the resolved docker executable '%s' could not be executed.",
|
||||
docker_exe, exc_info=True,
|
||||
error="Docker executable could not be executed. Check your Docker installation.",
|
||||
hint="Repair the Docker installation and retry.",
|
||||
)
|
||||
hint="Repair the Docker installation and retry.")
|
||||
except subprocess.TimeoutExpired:
|
||||
raise _docker_unavailable(
|
||||
"Docker backend selected but '%s version' timed out. The Docker daemon may not be running.",
|
||||
docker_exe, exc_info=True,
|
||||
error="Docker daemon is not responding. Ensure Docker is running and try again.",
|
||||
hint="Start the Docker daemon (e.g. `systemctl start docker` or "
|
||||
"launch Docker Desktop), then retry the same command.",
|
||||
)
|
||||
"launch Docker Desktop), then retry the same command.")
|
||||
except Exception:
|
||||
logger.error("Unexpected error while checking Docker availability.", exc_info=True)
|
||||
raise
|
||||
@@ -442,33 +397,28 @@ def _ensure_docker_available() -> None:
|
||||
docker_exe, result.returncode, result.stderr.strip(),
|
||||
error="Docker command is available but 'docker version' failed. Check your Docker installation.",
|
||||
hint="The Docker daemon may be down or the current user lacks "
|
||||
"permission (docker group). Fix and retry.",
|
||||
)
|
||||
"permission (docker group). Fix and retry.")
|
||||
|
||||
|
||||
def _name_only_env_args(names) -> list[str]:
|
||||
"""``-e KEY`` flags (no values): the docker CLI resolves values from its own env,
|
||||
so secrets live in owner-readable /proc/*/environ instead of world-readable cmdline."""
|
||||
args: list[str] = []
|
||||
for key in sorted(names):
|
||||
args.extend(["-e", key])
|
||||
return args
|
||||
return [arg for key in sorted(names) for arg in ("-e", key)]
|
||||
|
||||
|
||||
# Mount kinds declared by skills/credential_files: (getter, expects_file, log noun).
|
||||
_RO_MOUNT_SOURCES = (
|
||||
("get_credential_file_mounts", True, "credential"),
|
||||
("get_skills_directory_mount", False, "skills dir"),
|
||||
("get_cache_directory_mounts", False, "cache dir"),
|
||||
)
|
||||
("get_cache_directory_mounts", False, "cache dir"))
|
||||
|
||||
|
||||
def _readonly_skill_mount_args() -> list[str]:
|
||||
"""``-v host:container:ro`` args for credential files, skill dirs and cache dirs.
|
||||
|
||||
Read-only so the container can authenticate/read but never modify host state.
|
||||
Missing or wrong-kind sources are skipped with a warning (Docker-in-Docker
|
||||
auto-creates a missing file source as a directory, which would exit 125).
|
||||
Read-only so the container can authenticate/read but never modify host state. Missing
|
||||
or wrong-kind sources are skipped with a warning (Docker-in-Docker auto-creates a
|
||||
missing file source as a directory, which would exit 125).
|
||||
"""
|
||||
args: list[str] = []
|
||||
try:
|
||||
@@ -479,8 +429,7 @@ def _readonly_skill_mount_args() -> list[str]:
|
||||
if expects_file and src.is_dir():
|
||||
logger.warning(
|
||||
"Docker: skipping credential mount — source is a directory "
|
||||
"(likely Docker-in-Docker auto-creation): %s", src,
|
||||
)
|
||||
"(likely Docker-in-Docker auto-creation): %s", src)
|
||||
continue
|
||||
if expects_file and not src.is_file():
|
||||
logger.warning("Docker: skipping credential mount — source not found: %s", src)
|
||||
@@ -496,6 +445,22 @@ def _readonly_skill_mount_args() -> list[str]:
|
||||
return args
|
||||
|
||||
|
||||
def _host_user_args(run_as_host_user: bool) -> list[str]:
|
||||
"""``--user uid:gid`` so bind-mount writes are owned by the host user, not root. Without
|
||||
POSIX ids fall back to the full cap set — the image's init may still need to drop privileges."""
|
||||
if not run_as_host_user:
|
||||
return []
|
||||
user_spec = _resolve_host_user_spec()
|
||||
if user_spec is not None:
|
||||
logger.info("Docker: running container as host user %s", user_spec)
|
||||
return ["--user", user_spec]
|
||||
logger.warning(
|
||||
"docker_run_as_host_user is enabled but this platform does "
|
||||
"not expose POSIX uid/gid; container will start as its "
|
||||
"image default user.")
|
||||
return []
|
||||
|
||||
|
||||
class DockerEnvironment(BaseEnvironment):
|
||||
"""Hardened Docker container execution (caps dropped, no-new-privileges, PID limits,
|
||||
size-limited tmpfs). The container is the security boundary — its filesystem stays
|
||||
@@ -527,8 +492,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
extra_args: list = None,
|
||||
persist_across_processes: bool = True,
|
||||
shm_size: str = _DEFAULT_SHM_SIZE,
|
||||
shared_container_key: str = "",
|
||||
):
|
||||
shared_container_key: str = ""):
|
||||
if cwd == "~":
|
||||
cwd = "/root"
|
||||
super().__init__(cwd=cwd, timeout=timeout)
|
||||
@@ -542,11 +506,6 @@ class DockerEnvironment(BaseEnvironment):
|
||||
self._env = _normalize_env_dict(env)
|
||||
self._init_unset_passthrough_names: tuple[str, ...] = ()
|
||||
self._container_id: Optional[str] = None
|
||||
self._labels: dict[str, str] = {}
|
||||
self._image: str = ""
|
||||
self._image_uses_s6_init: bool = False
|
||||
self._all_run_args: list[str] = []
|
||||
self._run_env_values: dict[str, str] = {}
|
||||
self._init_env_values: dict[str, str] = {}
|
||||
self._workspace_dir: Optional[str] = None
|
||||
self._home_dir: Optional[str] = None
|
||||
@@ -560,40 +519,10 @@ class DockerEnvironment(BaseEnvironment):
|
||||
resource_args = self._resource_args(image, cpu, memory, disk, network, shm_size, extra_args)
|
||||
volume_args, writable_args = self._mount_args(volumes, host_cwd, auto_mount_cwd, task_id)
|
||||
volume_args.extend(_readonly_skill_mount_args())
|
||||
|
||||
# Egress credential-injection proxy: CA mount + HTTPS_PROXY/CA-bundle env
|
||||
# so outbound traffic routes through the host-side proxy and the sandbox
|
||||
# receives proxy tokens instead of real API keys.
|
||||
egress_volume_args, egress_env_overrides, egress_host_args = _egress_proxy_args_for_docker()
|
||||
egress_label = _egress_reuse_fingerprint(egress_volume_args, egress_env_overrides, egress_host_args)
|
||||
enforce_egress = _egress_enforce_on_docker() if egress_env_overrides else True
|
||||
critical_egress_names = _critical_egress_env_names(egress_env_overrides)
|
||||
if egress_env_overrides:
|
||||
check_forward_env_collisions(self._forward_env, critical_egress_names, enforce_egress)
|
||||
check_docker_env_collisions(self._env, egress_env_overrides, enforce_egress)
|
||||
egress_label, egress_volume_args, egress_host_args, env_args, validated_extra = (
|
||||
self._egress_and_env_args(extra_args))
|
||||
volume_args.extend(egress_volume_args)
|
||||
|
||||
merged_env = merge_egress_env(self._env, egress_env_overrides, enforce_egress)
|
||||
env_args = _name_only_env_args(merged_env)
|
||||
# Injected into the docker-client subprocess env at run time (also reused
|
||||
# verbatim by the container-recreation recovery path).
|
||||
self._run_env_values = dict(merged_env)
|
||||
|
||||
# Run as the host user so files written to bind mounts are owned by
|
||||
# them, not root. Without POSIX ids fall back to the full cap set —
|
||||
# the image's init may still need to drop privileges.
|
||||
user_args: list[str] = []
|
||||
if run_as_host_user:
|
||||
user_spec = _resolve_host_user_spec()
|
||||
if user_spec is not None:
|
||||
user_args = ["--user", user_spec]
|
||||
logger.info("Docker: running container as host user %s", user_spec)
|
||||
else:
|
||||
logger.warning(
|
||||
"docker_run_as_host_user is enabled but this platform does "
|
||||
"not expose POSIX uid/gid; container will start as its "
|
||||
"image default user."
|
||||
)
|
||||
user_args = _host_user_args(run_as_host_user)
|
||||
|
||||
# Resolved once so it works when /usr/local/bin is not in PATH (macOS services).
|
||||
self._docker_exe = find_docker() or "docker"
|
||||
@@ -603,25 +532,14 @@ class DockerEnvironment(BaseEnvironment):
|
||||
logger.info(
|
||||
"Docker: image %s uses /init (s6-overlay) as entrypoint — "
|
||||
"skipping --init and mounting /run with exec.",
|
||||
image,
|
||||
)
|
||||
image)
|
||||
security_args = _build_security_args(run_as_host_user and bool(user_args), run_exec=image_uses_s6_init)
|
||||
|
||||
logger.info("Docker volume_args: %s", volume_args)
|
||||
# docker_extra_args go last so they can override defaults.
|
||||
validated_extra = []
|
||||
for arg in (extra_args or []):
|
||||
if not isinstance(arg, str):
|
||||
logger.warning("Ignoring non-string docker_extra_args entry: %r", arg)
|
||||
continue
|
||||
validated_extra.append(arg)
|
||||
if egress_env_overrides:
|
||||
check_extra_args_collisions(validated_extra, critical_egress_names, enforce_egress)
|
||||
|
||||
all_run_args = (
|
||||
security_args + user_args + writable_args + resource_args
|
||||
+ egress_host_args + volume_args + env_args + validated_extra
|
||||
)
|
||||
+ egress_host_args + volume_args + env_args + validated_extra)
|
||||
logger.info("Docker run_args: %s", all_run_args)
|
||||
|
||||
# Labels identify hermes containers to the orphan reaper (hermes-agent=1),
|
||||
@@ -635,16 +553,14 @@ class DockerEnvironment(BaseEnvironment):
|
||||
"hermes-agent": "1",
|
||||
"hermes-task-id": task_label,
|
||||
"hermes-profile": profile_name,
|
||||
_EGRESS_LABEL_KEY: egress_label,
|
||||
}
|
||||
_EGRESS_LABEL_KEY: egress_label}
|
||||
# Saved for container recreation on "No such container" recovery.
|
||||
self._image = image
|
||||
self._image_uses_s6_init = image_uses_s6_init
|
||||
self._all_run_args = all_run_args
|
||||
|
||||
reused = persist_across_processes and self._attach_existing_container(
|
||||
task_label, profile_name, egress_label, network,
|
||||
)
|
||||
task_label, profile_name, egress_label, network)
|
||||
if not reused:
|
||||
self._container_id = self._docker_run(cwd)
|
||||
|
||||
@@ -656,6 +572,34 @@ class DockerEnvironment(BaseEnvironment):
|
||||
# __init__ helpers
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
def _egress_and_env_args(self, extra_args) -> tuple[str, list[str], list[str], list[str], list[str]]:
|
||||
"""Egress credential-injection proxy plumbing (CA mount + HTTPS_PROXY/CA-bundle env so
|
||||
outbound traffic routes through the host-side proxy and the sandbox receives proxy tokens
|
||||
instead of real API keys), merged with docker_env into name-only ``-e`` args, plus the
|
||||
validated docker_extra_args. Returns ``(egress_label, volume_args, host_args, env_args,
|
||||
validated_extra)``; sets ``self._run_env_values`` (injected into the docker-client
|
||||
subprocess env at run time and reused verbatim by container-recreation recovery)."""
|
||||
egress_volume_args, egress_env_overrides, egress_host_args = _egress_proxy_args_for_docker()
|
||||
egress_label = _egress_reuse_fingerprint(egress_volume_args, egress_env_overrides, egress_host_args)
|
||||
enforce_egress = _egress_enforce_on_docker() if egress_env_overrides else True
|
||||
critical_egress_names = _critical_egress_env_names(egress_env_overrides)
|
||||
if egress_env_overrides:
|
||||
check_forward_env_collisions(self._forward_env, critical_egress_names, enforce_egress)
|
||||
check_docker_env_collisions(self._env, egress_env_overrides, enforce_egress)
|
||||
|
||||
merged_env = merge_egress_env(self._env, egress_env_overrides, enforce_egress)
|
||||
self._run_env_values = dict(merged_env)
|
||||
|
||||
validated_extra = []
|
||||
for arg in (extra_args or []):
|
||||
if not isinstance(arg, str):
|
||||
logger.warning("Ignoring non-string docker_extra_args entry: %r", arg)
|
||||
continue
|
||||
validated_extra.append(arg)
|
||||
if egress_env_overrides:
|
||||
check_extra_args_collisions(validated_extra, critical_egress_names, enforce_egress)
|
||||
return egress_label, egress_volume_args, egress_host_args, _name_only_env_args(merged_env), validated_extra
|
||||
|
||||
def _resource_args(self, image, cpu, memory, disk, network, shm_size, extra_args) -> list[str]:
|
||||
"""cgroup-gated CPU/memory/pids limits, shm size, disk quota and network mode."""
|
||||
args: list[str] = []
|
||||
@@ -676,8 +620,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
else:
|
||||
logger.warning(
|
||||
"Docker storage driver does not support per-container disk limits "
|
||||
"(requires overlay2 on XFS with pquota). Container will run without disk quota."
|
||||
)
|
||||
"(requires overlay2 on XFS with pquota). Container will run without disk quota.")
|
||||
if not network:
|
||||
args.append("--network=none")
|
||||
return args
|
||||
@@ -704,8 +647,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
host_cwd_abs = os.path.abspath(os.path.expanduser(host_cwd)) if host_cwd else ""
|
||||
bind_host_cwd = (
|
||||
auto_mount_cwd and bool(host_cwd_abs) and os.path.isdir(host_cwd_abs)
|
||||
and not workspace_explicitly_mounted
|
||||
)
|
||||
and not workspace_explicitly_mounted)
|
||||
if auto_mount_cwd and host_cwd and not os.path.isdir(host_cwd_abs):
|
||||
logger.debug("Skipping docker cwd mount: host_cwd is not a valid directory: %s", host_cwd)
|
||||
mount_workspace = not bind_host_cwd and not workspace_explicitly_mounted
|
||||
@@ -753,8 +695,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
"docker_network=false requests an air-gapped "
|
||||
"container — removing it and starting fresh "
|
||||
"(task=%s, profile=%s).",
|
||||
container_id[:12], actual_mode or "unknown", task_label, profile_name,
|
||||
)
|
||||
container_id[:12], actual_mode or "unknown", task_label, profile_name)
|
||||
try:
|
||||
run_capture([self._docker_exe, "rm", "-f", container_id], timeout=30)
|
||||
except (subprocess.TimeoutExpired, OSError) as e:
|
||||
@@ -768,22 +709,18 @@ class DockerEnvironment(BaseEnvironment):
|
||||
logger.warning(
|
||||
"Failed to start existing container %s (state=%s): "
|
||||
"%s — falling back to a fresh container.",
|
||||
container_id[:12], state, e,
|
||||
)
|
||||
container_id[:12], state, e)
|
||||
return False
|
||||
self._container_id = container_id
|
||||
logger.info(
|
||||
"Reusing container %s (task=%s, profile=%s, prior state=%s)",
|
||||
container_id[:12], task_label, profile_name, state,
|
||||
)
|
||||
container_id[:12], task_label, profile_name, state)
|
||||
return True
|
||||
|
||||
def _run_command(self, name: str, workdir: str) -> list[str]:
|
||||
"""``docker run -d`` argv for a fresh ``sleep infinity`` container (idle reaper handles
|
||||
lifetime). s6-overlay images already provide PID 1, so ``--init`` is skipped for them."""
|
||||
label_args: list[str] = []
|
||||
for k, v in self._labels.items():
|
||||
label_args.extend(["--label", f"{k}={v}"])
|
||||
label_args = [arg for k, v in self._labels.items() for arg in ("--label", f"{k}={v}")]
|
||||
return [
|
||||
self._docker_exe, "run", "-d",
|
||||
*([] if self._image_uses_s6_init else ["--init"]),
|
||||
@@ -792,8 +729,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
"-w", workdir,
|
||||
*self._all_run_args,
|
||||
self._image,
|
||||
"sleep", "infinity",
|
||||
]
|
||||
"sleep", "infinity"]
|
||||
|
||||
def _docker_run(self, cwd: str) -> str:
|
||||
"""Start a fresh container and return its id. A failed ``docker run`` (exit 125, timeout
|
||||
@@ -805,14 +741,12 @@ class DockerEnvironment(BaseEnvironment):
|
||||
try:
|
||||
result = run_capture(
|
||||
run_cmd, timeout=120, check=True, # image pull may take a while
|
||||
env=self._docker_client_env(self._run_env_values),
|
||||
)
|
||||
env=self._docker_client_env(self._run_env_values))
|
||||
except (subprocess.CalledProcessError, subprocess.TimeoutExpired) as e:
|
||||
logger.warning("docker run failed for %s, cleaning up orphaned container: %s", container_name, e)
|
||||
subprocess.run(
|
||||
[self._docker_exe, "rm", "-f", container_name],
|
||||
capture_output=True, timeout=10, stdin=subprocess.DEVNULL,
|
||||
)
|
||||
capture_output=True, timeout=10, stdin=subprocess.DEVNULL)
|
||||
raise
|
||||
container_id = result.stdout.strip()
|
||||
logger.info("Started container %s (%s)", container_name, container_id[:12])
|
||||
@@ -825,9 +759,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
def _docker_client_env(self, values: dict[str, str]) -> dict[str, str] | None:
|
||||
"""Env for the docker-client subprocess carrying forwarded values (pairs with name-only
|
||||
``-e KEY`` flags to keep secrets out of cmdline); ``None`` = inherit when empty."""
|
||||
if not values:
|
||||
return None
|
||||
return {**os.environ, **values}
|
||||
return {**os.environ, **values} if values else None
|
||||
|
||||
def _build_init_env_args(self) -> list[str]:
|
||||
"""Name-only ``-e`` args for init_session so ``export -p`` captures docker_env
|
||||
@@ -928,8 +860,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
existing = self._find_reusable_container(
|
||||
self._labels.get("hermes-task-id", ""),
|
||||
self._labels.get("hermes-profile", ""),
|
||||
self._labels.get(_EGRESS_LABEL_KEY, "off"),
|
||||
)
|
||||
self._labels.get(_EGRESS_LABEL_KEY, "off"))
|
||||
if existing is not None:
|
||||
cid, state = existing
|
||||
if state == "running":
|
||||
@@ -951,8 +882,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
new_name = f"hermes-{uuid.uuid4().hex[:8]}"
|
||||
result = run_capture(
|
||||
self._run_command(new_name, self.cwd), timeout=120, check=True,
|
||||
env=self._docker_client_env(self._run_env_values),
|
||||
)
|
||||
env=self._docker_client_env(self._run_env_values))
|
||||
self._container_id = result.stdout.strip()
|
||||
logger.info("Recovery: created fresh container %s (%s)", new_name, self._container_id[:12])
|
||||
except (subprocess.CalledProcessError, subprocess.TimeoutExpired, OSError) as e:
|
||||
@@ -977,8 +907,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
result.get("returncode", 0) != 0
|
||||
and self._is_container_gone(result.get("output", ""))
|
||||
and self._persist_across_processes
|
||||
and self._recreate_container()
|
||||
):
|
||||
and self._recreate_container()):
|
||||
result = super().execute(command, cwd, **kwargs)
|
||||
return result
|
||||
|
||||
@@ -1011,8 +940,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
try:
|
||||
result = run_capture(
|
||||
[self._docker_exe, "inspect", "--format", "{{.HostConfig.NetworkMode}}", container_id],
|
||||
timeout=10,
|
||||
)
|
||||
timeout=10)
|
||||
except (subprocess.TimeoutExpired, OSError) as e:
|
||||
logger.debug("docker inspect NetworkMode failed: %s", e)
|
||||
return None
|
||||
@@ -1022,11 +950,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
return result.stdout.strip() or None
|
||||
|
||||
def _find_reusable_container(
|
||||
self,
|
||||
task_label: str,
|
||||
profile_label: str,
|
||||
egress_label: str,
|
||||
) -> Optional[tuple[str, str]]:
|
||||
self, task_label: str, profile_label: str, egress_label: str) -> Optional[tuple[str, str]]:
|
||||
"""``(container_id, state)`` of an existing container labeled for this task/profile/
|
||||
egress posture, or ``None`` on miss or any failure. With egress off the probe is
|
||||
widened to all task+profile containers and post-filtered to reject a non-"off" egress
|
||||
@@ -1036,8 +960,7 @@ class DockerEnvironment(BaseEnvironment):
|
||||
filters = [
|
||||
"--filter", "label=hermes-agent=1",
|
||||
"--filter", f"label=hermes-task-id={task_label}",
|
||||
"--filter", f"label=hermes-profile={profile_label}",
|
||||
]
|
||||
"--filter", f"label=hermes-profile={profile_label}"]
|
||||
if egress_off:
|
||||
fmt = '{{.ID}}\t{{.State}}\t{{.Label "' + _EGRESS_LABEL_KEY + '"}}'
|
||||
else:
|
||||
@@ -1051,22 +974,20 @@ class DockerEnvironment(BaseEnvironment):
|
||||
if result.returncode != 0:
|
||||
logger.debug(
|
||||
"docker ps probe returned %d: %s — will start a fresh container",
|
||||
result.returncode, result.stderr.strip(),
|
||||
)
|
||||
result.returncode, result.stderr.strip())
|
||||
return None
|
||||
# Multiple matches can happen after a crash mid-cleanup: prefer a running
|
||||
# one, else the first listed; stale duplicates are the orphan reaper's job.
|
||||
running = None
|
||||
first = None
|
||||
nparts = 3 if egress_off else 2
|
||||
running = first = None
|
||||
for ln in (ln for ln in result.stdout.splitlines() if ln.strip()):
|
||||
parts = ln.split("\t", 2 if egress_off else 1)
|
||||
if len(parts) != (3 if egress_off else 2):
|
||||
parts = ln.split("\t", nparts - 1)
|
||||
if len(parts) != nparts:
|
||||
continue
|
||||
cid, state = parts[0], parts[1].lower()
|
||||
if egress_off and parts[2] not in ("", "<no value>", "off"):
|
||||
logger.debug(
|
||||
"skipping container %s for egress=off reuse: label %s=%r", cid, _EGRESS_LABEL_KEY, parts[2],
|
||||
)
|
||||
"skipping container %s for egress=off reuse: label %s=%r", cid, _EGRESS_LABEL_KEY, parts[2])
|
||||
continue
|
||||
if first is None:
|
||||
first = (cid, state)
|
||||
@@ -1104,17 +1025,14 @@ class DockerEnvironment(BaseEnvironment):
|
||||
log_id = container_id[:12]
|
||||
|
||||
def _do_cleanup() -> None:
|
||||
for verb, argv in (("stop", ["stop", "-t", "10"]), ("rm", ["rm", "-f"])):
|
||||
for argv, fail_msg in ((["stop", "-t", "10"], "docker stop %s timed out / failed: %s"),
|
||||
(["rm", "-f"], "docker rm -f %s failed: %s")):
|
||||
try:
|
||||
subprocess.run(
|
||||
[docker_exe, *argv, container_id],
|
||||
capture_output=True, timeout=30, stdin=subprocess.DEVNULL,
|
||||
)
|
||||
capture_output=True, timeout=30, stdin=subprocess.DEVNULL)
|
||||
except (subprocess.TimeoutExpired, OSError) as e:
|
||||
if verb == "stop":
|
||||
logger.warning("docker stop %s timed out / failed: %s", log_id, e)
|
||||
else:
|
||||
logger.warning("docker rm -f %s failed: %s", log_id, e)
|
||||
logger.warning(fail_msg, log_id, e)
|
||||
|
||||
t = threading.Thread(target=_do_cleanup, daemon=True, name=f"hermes-cleanup-{log_id}")
|
||||
t.start()
|
||||
|
||||
@@ -23,8 +23,7 @@ _PROXY_CONTROL_ENV = frozenset({
|
||||
"HTTPS_PROXY", "https_proxy", "HTTP_PROXY", "http_proxy",
|
||||
"NO_PROXY", "no_proxy",
|
||||
"REQUESTS_CA_BUNDLE", "SSL_CERT_FILE", "CURL_CA_BUNDLE",
|
||||
"NODE_EXTRA_CA_CERTS",
|
||||
})
|
||||
"NODE_EXTRA_CA_CERTS"})
|
||||
|
||||
|
||||
def _egress_proxy_args_for_docker() -> tuple[list[str], dict[str, str], list[str]]:
|
||||
@@ -56,20 +55,17 @@ def _egress_proxy_args_for_docker() -> tuple[list[str], dict[str, str], list[str
|
||||
if not status.configured:
|
||||
return _degraded(
|
||||
"proxy.enabled is true but iron-proxy is not configured. "
|
||||
"Run `hermes egress setup` to mint tokens and write proxy.yaml."
|
||||
)
|
||||
"Run `hermes egress setup` to mint tokens and write proxy.yaml.")
|
||||
if not (status.pid and status.listening):
|
||||
return _degraded(
|
||||
f"iron-proxy is enabled but not running on port {status.tunnel_port}. "
|
||||
"Start it with `hermes egress start`."
|
||||
)
|
||||
"Start it with `hermes egress start`.")
|
||||
if status.ca_cert_path is None or not status.ca_cert_path.exists():
|
||||
# Configured a moment ago but the trust anchor vanished: proxy env vars
|
||||
# without the CA would make every TLS handshake fail.
|
||||
return _degraded(
|
||||
f"iron-proxy CA cert vanished from {status.ca_cert_path}. "
|
||||
"Re-run `hermes egress setup` to regenerate it."
|
||||
)
|
||||
"Re-run `hermes egress setup` to regenerate it.")
|
||||
# Empty/corrupt mappings look like an upstream outage from inside the
|
||||
# sandbox (every request 403s); refuse rather than ship a broken sandbox.
|
||||
mappings = ip.load_mappings()
|
||||
@@ -77,8 +73,7 @@ def _egress_proxy_args_for_docker() -> tuple[list[str], dict[str, str], list[str
|
||||
return _degraded(
|
||||
"iron-proxy is configured but mappings.json is empty or "
|
||||
"corrupt. Re-run `hermes egress setup` to mint provider "
|
||||
"tokens before starting a sandbox."
|
||||
)
|
||||
"tokens before starting a sandbox.")
|
||||
|
||||
volume_args = ["-v", f"{status.ca_cert_path}:{_CONTAINER_CA}:ro"]
|
||||
|
||||
@@ -103,8 +98,7 @@ def _egress_proxy_args_for_docker() -> tuple[list[str], dict[str, str], list[str
|
||||
"CURL_CA_BUNDLE": _CONTAINER_CA,
|
||||
"NODE_EXTRA_CA_CERTS": _CONTAINER_CA,
|
||||
"HERMES_EGRESS_PROXY": "1", # lets the in-sandbox agent know it is proxy-aware
|
||||
_NODE_OPTIONS_SENTINEL: "--use-openssl-ca",
|
||||
}
|
||||
_NODE_OPTIONS_SENTINEL: "--use-openssl-ca"}
|
||||
|
||||
# Proxy tokens under the standard provider env names (and their aliases) so
|
||||
# SDKs work unchanged; HERMES_PROXY_TOKEN_* copies are for diagnostics.
|
||||
@@ -120,15 +114,13 @@ def _egress_proxy_args_for_docker() -> tuple[list[str], dict[str, str], list[str
|
||||
|
||||
|
||||
def _egress_reuse_fingerprint(
|
||||
volume_args: list[str], env_overrides: dict[str, str], host_args: list[str],
|
||||
) -> str:
|
||||
volume_args: list[str], env_overrides: dict[str, str], host_args: list[str]) -> str:
|
||||
"""Stable Docker-label value for the egress posture of a container."""
|
||||
if not (volume_args or env_overrides or host_args):
|
||||
return "off"
|
||||
payload = json.dumps(
|
||||
{"volume_args": volume_args, "env_overrides": env_overrides, "host_args": host_args},
|
||||
sort_keys=True, separators=(",", ":"),
|
||||
)
|
||||
sort_keys=True, separators=(",", ":"))
|
||||
return hashlib.sha256(payload.encode("utf-8")).hexdigest()[:24]
|
||||
|
||||
|
||||
@@ -163,10 +155,8 @@ def _extra_args_egress_collisions(extra_args: list[str], critical_names: set[str
|
||||
if arg in env_flags:
|
||||
if arg == "--env-file":
|
||||
collisions.append(arg)
|
||||
else:
|
||||
name = nxt.split("=", 1)[0]
|
||||
if name in critical_names:
|
||||
collisions.append(name)
|
||||
elif nxt.split("=", 1)[0] in critical_names:
|
||||
collisions.append(nxt.split("=", 1)[0])
|
||||
i += 2
|
||||
continue
|
||||
if any(arg.startswith(f"{flag}=") for flag in env_flags):
|
||||
@@ -198,8 +188,7 @@ def check_forward_env_collisions(forward_env: list[str], critical: set[str], enf
|
||||
enforce=enforce,
|
||||
remedy="Remove these names from docker_forward_env or disable enforce_on_docker "
|
||||
"to opt out of egress isolation.",
|
||||
consequence="Explicit docker_forward_env values will override egress tokens.",
|
||||
)
|
||||
consequence="Explicit docker_forward_env values will override egress tokens.")
|
||||
|
||||
|
||||
def check_docker_env_collisions(user_env: dict[str, str], egress_env: dict[str, str], enforce: bool) -> None:
|
||||
@@ -211,6 +200,7 @@ def check_docker_env_collisions(user_env: dict[str, str], egress_env: dict[str,
|
||||
provider_keys = {m.real_env_name for m in ip.load_mappings()}
|
||||
except Exception: # best-effort
|
||||
pass
|
||||
|
||||
def _collides(k: str) -> bool:
|
||||
if k not in user_env:
|
||||
return False
|
||||
@@ -226,8 +216,7 @@ def check_docker_env_collisions(user_env: dict[str, str], egress_env: dict[str,
|
||||
remedy="Remove these keys from docker_env or disable enforce_on_docker to opt out "
|
||||
"of egress isolation.",
|
||||
consequence="Falling back to docker_env values; sandbox traffic will NOT route "
|
||||
"through the proxy.",
|
||||
)
|
||||
"through the proxy.")
|
||||
|
||||
|
||||
def check_extra_args_collisions(extra_args: list[str], critical: set[str], enforce: bool) -> None:
|
||||
@@ -237,18 +226,14 @@ def check_extra_args_collisions(extra_args: list[str], critical: set[str], enfor
|
||||
f"docker_extra_args would override egress-proxy controls {collisions}",
|
||||
enforce=enforce,
|
||||
remedy="Remove these args or disable enforce_on_docker to opt out of egress isolation.",
|
||||
consequence="Extra Docker args may bypass egress isolation.",
|
||||
)
|
||||
consequence="Extra Docker args may bypass egress isolation.")
|
||||
|
||||
|
||||
def merge_egress_env(user_env: dict[str, str], egress_env: dict[str, str], enforce: bool) -> dict[str, str]:
|
||||
"""Merge docker_env with egress overrides (egress wins under enforcement, docker_env
|
||||
otherwise) and resolve the NODE_OPTIONS sentinel: the flag is APPENDED to the operator's
|
||||
NODE_OPTIONS after stripping conflicting CA-mode flags so it wins deterministically."""
|
||||
if enforce and egress_env:
|
||||
merged = {**user_env, **egress_env}
|
||||
else:
|
||||
merged = {**egress_env, **user_env}
|
||||
merged = {**user_env, **egress_env} if enforce and egress_env else {**egress_env, **user_env}
|
||||
|
||||
raw_append = merged.pop(_NODE_OPTIONS_SENTINEL, None)
|
||||
if raw_append:
|
||||
@@ -260,8 +245,7 @@ def merge_egress_env(user_env: dict[str, str], egress_env: dict[str, str], enfor
|
||||
logger.warning(
|
||||
"Overriding conflicting NODE_OPTIONS CA-mode flag(s) %s "
|
||||
"with egress-required %s to keep Node routed through the "
|
||||
"egress CA store.", dropped, append_token,
|
||||
)
|
||||
"egress CA store.", dropped, append_token)
|
||||
tokens = [t for t in tokens if t not in _CA_MODE_FLAGS or t == append_token]
|
||||
if append_token not in tokens:
|
||||
tokens.append(append_token)
|
||||
|
||||
+44
-103
@@ -51,25 +51,18 @@ _SYNC_BACK_MAX_BYTES = 2 * 1024 * 1024 * 1024 # 2 GiB — refuse to extract lar
|
||||
|
||||
|
||||
def iter_sync_files(container_base: str = "/root/.hermes") -> list[tuple[str, str]]:
|
||||
"""Enumerate all (host_path, remote_path) pairs to sync to a remote.
|
||||
|
||||
Credential paths are remapped from the hardcoded /root/.hermes to
|
||||
*container_base* because the remote user's home may differ.
|
||||
"""
|
||||
"""Enumerate all (host_path, remote_path) pairs to sync to a remote. Credential paths are
|
||||
remapped from the hardcoded /root/.hermes to *container_base* (remote home may differ)."""
|
||||
# Late import: credential_files pulls in agent modules (circular at module level).
|
||||
from tools.credential_files import (
|
||||
get_credential_file_mounts,
|
||||
iter_cache_files,
|
||||
iter_skills_files,
|
||||
)
|
||||
from tools.credential_files import get_credential_file_mounts, iter_cache_files, iter_skills_files
|
||||
|
||||
files: list[tuple[str, str]] = [
|
||||
files = [
|
||||
(entry["host_path"], entry["container_path"].replace("/root/.hermes", container_base, 1))
|
||||
for entry in get_credential_file_mounts()
|
||||
]
|
||||
for entry in (*iter_skills_files(container_base=container_base),
|
||||
*iter_cache_files(container_base=container_base)):
|
||||
files.append((entry["host_path"], entry["container_path"]))
|
||||
for entry in get_credential_file_mounts()]
|
||||
files += [
|
||||
(entry["host_path"], entry["container_path"])
|
||||
for entry in (*iter_skills_files(container_base=container_base),
|
||||
*iter_cache_files(container_base=container_base))]
|
||||
return files
|
||||
|
||||
|
||||
@@ -85,15 +78,12 @@ def _credential_host_paths() -> set[str]:
|
||||
"""Return credential files that are upload-only for remote sandboxes."""
|
||||
try:
|
||||
from tools.credential_files import get_credential_file_mounts
|
||||
|
||||
mounts = get_credential_file_mounts()
|
||||
except Exception:
|
||||
return set()
|
||||
return {
|
||||
_resolve_host_path_str(entry["host_path"])
|
||||
for entry in mounts
|
||||
if isinstance(entry, dict) and entry.get("host_path")
|
||||
}
|
||||
for entry in mounts if isinstance(entry, dict) and entry.get("host_path")}
|
||||
|
||||
|
||||
def quoted_rm_command(remote_paths: list[str]) -> str:
|
||||
@@ -121,12 +111,9 @@ def _sha256_file(path: str) -> str:
|
||||
|
||||
|
||||
class FileSyncManager:
|
||||
"""Tracks local file changes and syncs to a remote environment.
|
||||
|
||||
Backends instantiate this with transport callbacks (upload, delete) and a
|
||||
file-source callable. The manager handles mtime-based change detection,
|
||||
deletion tracking, rate limiting, and transactional state.
|
||||
"""
|
||||
"""Tracks local file changes and syncs to a remote environment. Backends supply transport
|
||||
callbacks (upload, delete) and a file-source callable; the manager handles mtime-based
|
||||
change detection, deletion tracking, rate limiting, and transactional state."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
@@ -135,8 +122,7 @@ class FileSyncManager:
|
||||
delete_fn: DeleteFn,
|
||||
sync_interval: float = _SYNC_INTERVAL_SECONDS,
|
||||
bulk_upload_fn: BulkUploadFn | None = None,
|
||||
bulk_download_fn: BulkDownloadFn | None = None,
|
||||
):
|
||||
bulk_download_fn: BulkDownloadFn | None = None):
|
||||
self._get_files_fn = get_files_fn
|
||||
self._upload_fn = upload_fn
|
||||
self._bulk_upload_fn = bulk_upload_fn
|
||||
@@ -150,13 +136,10 @@ class FileSyncManager:
|
||||
self._sync_interval = sync_interval
|
||||
|
||||
def sync(self, *, force: bool = False) -> None:
|
||||
"""Run a sync cycle: upload changed files, delete removed files.
|
||||
|
||||
Rate-limited to once per ``sync_interval`` unless *force* is True or
|
||||
``HERMES_FORCE_FILE_SYNC=1`` is set. Transactional: state is committed
|
||||
only if ALL operations succeed; on failure it rolls back so the next
|
||||
cycle retries everything.
|
||||
"""
|
||||
"""Run a sync cycle: upload changed files, delete removed files. Rate-limited to once
|
||||
per ``sync_interval`` unless *force* or ``HERMES_FORCE_FILE_SYNC=1``. Transactional:
|
||||
state is committed only if ALL operations succeed; on failure it rolls back so the
|
||||
next cycle retries everything."""
|
||||
with self._transaction_lock:
|
||||
self._sync_transaction(force=force)
|
||||
|
||||
@@ -165,8 +148,7 @@ class FileSyncManager:
|
||||
if (
|
||||
not force
|
||||
and not os.environ.get(_FORCE_SYNC_ENV)
|
||||
and _monotonic() - self._last_sync_time < self._sync_interval
|
||||
):
|
||||
and _monotonic() - self._last_sync_time < self._sync_interval):
|
||||
return
|
||||
|
||||
current_files = self._get_files_fn()
|
||||
@@ -192,20 +174,14 @@ class FileSyncManager:
|
||||
except Exception as exc:
|
||||
self._synced_files = prev_files
|
||||
self._pushed_hashes = prev_hashes
|
||||
# Do NOT advance _last_sync_time: bumping the rate-limit clock on
|
||||
# failure would suppress the retry for up to _sync_interval,
|
||||
# contradicting the "next cycle retries everything" contract.
|
||||
# Do NOT advance _last_sync_time: bumping the rate-limit clock on failure would
|
||||
# suppress the retry for up to _sync_interval, contradicting the retry contract.
|
||||
logger.warning("file_sync: sync failed, rolled back state: %s", exc)
|
||||
|
||||
def _plan_sync(
|
||||
self, current_files: list[tuple[str, str]]
|
||||
) -> tuple[list[tuple[str, str]], dict[str, tuple[float, int]], list[str]]:
|
||||
"""Diff *current_files* against synced state.
|
||||
|
||||
Returns ``(to_upload, new_synced_state, to_delete)``: new/changed
|
||||
(mtime,size) pairs to upload, the state to commit if everything
|
||||
succeeds, and synced remote paths no longer present locally.
|
||||
"""
|
||||
"""Diff *current_files* against synced state -> ``(to_upload, new_synced_state, to_delete)``."""
|
||||
to_upload: list[tuple[str, str]] = []
|
||||
new_files = dict(self._synced_files)
|
||||
for host_path, remote_path in current_files:
|
||||
@@ -240,13 +216,9 @@ class FileSyncManager:
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
def sync_back(self, hermes_home: Path | None = None) -> None:
|
||||
"""Pull remote changes back to the host filesystem.
|
||||
|
||||
Downloads the remote ``.hermes/`` directory as a tar, unpacks it, and
|
||||
applies only files whose SHA-256 differs from what was pushed. SIGINT is
|
||||
deferred until complete; concurrent gateway sandboxes are serialized
|
||||
via a file lock.
|
||||
"""
|
||||
"""Pull remote changes back to the host: download the remote ``.hermes/`` as a tar and
|
||||
apply only files whose SHA-256 differs from what was pushed. SIGINT is deferred until
|
||||
complete; concurrent gateway sandboxes are serialized via a file lock."""
|
||||
with self._transaction_lock:
|
||||
self._sync_back_transaction(hermes_home=hermes_home)
|
||||
|
||||
@@ -273,10 +245,7 @@ class FileSyncManager:
|
||||
last_exc = exc
|
||||
if attempt < _SYNC_BACK_MAX_RETRIES - 1:
|
||||
delay = _SYNC_BACK_BACKOFF[attempt]
|
||||
logger.warning(
|
||||
"sync_back: attempt %d failed (%s), retrying in %ds",
|
||||
attempt + 1, exc, delay,
|
||||
)
|
||||
logger.warning("sync_back: attempt %d failed (%s), retrying in %ds", attempt + 1, exc, delay)
|
||||
_sleep(delay)
|
||||
|
||||
logger.warning("sync_back: all %d attempts failed: %s", _SYNC_BACK_MAX_RETRIES, last_exc)
|
||||
@@ -303,10 +272,9 @@ class FileSyncManager:
|
||||
if on_main_thread and original_handler is not None:
|
||||
signal.signal(signal.SIGINT, original_handler)
|
||||
if deferred_sigint:
|
||||
# Re-deliver the deferred Ctrl+C to the restored handler.
|
||||
# ``os.kill(os.getpid(), SIGINT)`` is NOT graceful on
|
||||
# Windows (routes to TerminateProcess, hard-killing the
|
||||
# CLI); ``raise_signal`` invokes the handler everywhere.
|
||||
# Re-deliver the deferred Ctrl+C to the restored handler. ``os.kill(os.getpid(),
|
||||
# SIGINT)`` is NOT graceful on Windows (routes to TerminateProcess, hard-killing
|
||||
# the CLI); ``raise_signal`` invokes the handler everywhere.
|
||||
signal.raise_signal(signal.SIGINT)
|
||||
|
||||
def _sync_back_locked(self, lock_path: Path) -> None:
|
||||
@@ -348,8 +316,7 @@ class FileSyncManager:
|
||||
if tar_size > _SYNC_BACK_MAX_BYTES:
|
||||
logger.warning(
|
||||
"sync_back: remote tar is %d bytes (cap %d) — skipping extraction",
|
||||
tar_size, _SYNC_BACK_MAX_BYTES,
|
||||
)
|
||||
tar_size, _SYNC_BACK_MAX_BYTES)
|
||||
return
|
||||
|
||||
with tempfile.TemporaryDirectory(prefix="hermes-sync-back-") as staging:
|
||||
@@ -362,9 +329,7 @@ class FileSyncManager:
|
||||
for fname in filenames:
|
||||
staged_file = os.path.join(dirpath, fname)
|
||||
remote_path = "/" + os.path.relpath(staged_file, staging)
|
||||
applied += self._apply_staged_file(
|
||||
staged_file, remote_path, file_mapping, upload_only
|
||||
)
|
||||
applied += self._apply_staged_file(staged_file, remote_path, file_mapping, upload_only)
|
||||
|
||||
if applied:
|
||||
logger.info("sync_back: applied %d changed file(s)", applied)
|
||||
@@ -372,27 +337,18 @@ class FileSyncManager:
|
||||
logger.debug("sync_back: no remote changes detected")
|
||||
|
||||
def _apply_staged_file(
|
||||
self,
|
||||
staged_file: str,
|
||||
remote_path: str,
|
||||
file_mapping: list[tuple[str, str]],
|
||||
upload_only_host_paths: set[str],
|
||||
self, staged_file: str, remote_path: str, file_mapping: list[tuple[str, str]], upload_only_host_paths: set[str],
|
||||
) -> int:
|
||||
"""Copy one extracted remote file onto the host if it changed since push.
|
||||
|
||||
Returns 1 if applied, 0 if skipped (unchanged, unmapped, or an
|
||||
upload-only credential). A host file modified since push is
|
||||
overwritten with the remote version (last-write-wins) with a warning.
|
||||
"""
|
||||
"""Copy one extracted remote file onto the host if it changed since push. Returns 1 if
|
||||
applied, 0 if skipped (unchanged, unmapped, or an upload-only credential). A host file
|
||||
modified since push is overwritten with the remote version (last-write-wins) with a warning."""
|
||||
pushed_hash = self._pushed_hashes.get(remote_path)
|
||||
if pushed_hash is not None and _sha256_file(staged_file) == pushed_hash:
|
||||
return 0 # unchanged from push
|
||||
|
||||
host_path = self._resolve_host_path(remote_path, file_mapping)
|
||||
if host_path is None:
|
||||
host_path = self._infer_host_path(
|
||||
remote_path, file_mapping, upload_only_host_paths=upload_only_host_paths
|
||||
)
|
||||
host_path = self._infer_host_path(remote_path, file_mapping, upload_only_host_paths=upload_only_host_paths)
|
||||
if host_path is None:
|
||||
logger.debug("sync_back: skipping %s (no host mapping)", remote_path)
|
||||
return 0
|
||||
@@ -401,41 +357,26 @@ class FileSyncManager:
|
||||
logger.debug("sync_back: skipping upload-only credential file %s", remote_path)
|
||||
return 0
|
||||
|
||||
if (
|
||||
pushed_hash is not None
|
||||
and os.path.exists(host_path)
|
||||
and _sha256_file(host_path) != pushed_hash
|
||||
):
|
||||
if pushed_hash is not None and os.path.exists(host_path) and _sha256_file(host_path) != pushed_hash:
|
||||
logger.warning(
|
||||
"sync_back: conflict on %s — host modified "
|
||||
"since push, remote also changed. Applying "
|
||||
"remote version (last-write-wins).",
|
||||
remote_path,
|
||||
)
|
||||
remote_path)
|
||||
|
||||
os.makedirs(os.path.dirname(host_path), exist_ok=True)
|
||||
shutil.copy2(staged_file, host_path)
|
||||
return 1
|
||||
|
||||
def _resolve_host_path(self, remote_path: str,
|
||||
file_mapping: list[tuple[str, str]] | None = None) -> str | None:
|
||||
def _resolve_host_path(self, remote_path: str, file_mapping: list[tuple[str, str]] | None = None) -> str | None:
|
||||
"""Find the host path for a known remote path from the file mapping."""
|
||||
for host, remote in file_mapping or []:
|
||||
if remote == remote_path:
|
||||
return host
|
||||
return None
|
||||
return next((host for host, remote in file_mapping or [] if remote == remote_path), None)
|
||||
|
||||
def _infer_host_path(self, remote_path: str,
|
||||
file_mapping: list[tuple[str, str]] | None = None,
|
||||
*,
|
||||
def _infer_host_path(self, remote_path: str, file_mapping: list[tuple[str, str]] | None = None, *,
|
||||
upload_only_host_paths: set[str] | None = None) -> str | None:
|
||||
"""Infer a host path for a new remote file by matching path prefixes.
|
||||
|
||||
Uses an existing remote->host pair whose parent directory prefixes
|
||||
*remote_path* and applies the same substitution, e.g. mapping
|
||||
``/root/.hermes/skills/a.md`` -> ``~/.hermes/skills/a.md`` sends a new
|
||||
``/root/.hermes/skills/b.md`` to ``~/.hermes/skills/b.md``.
|
||||
"""
|
||||
"""Infer a host path for a new remote file by matching path prefixes: an existing
|
||||
remote->host pair whose parent directory prefixes *remote_path* gets the same
|
||||
substitution (``/root/.hermes/skills/b.md`` -> ``~/.hermes/skills/b.md``)."""
|
||||
upload_only_host_paths = upload_only_host_paths or set()
|
||||
for host, remote in file_mapping or []:
|
||||
if self._is_upload_only_host_path(host, upload_only_host_paths):
|
||||
|
||||
@@ -1,9 +1,6 @@
|
||||
"""Path-component helpers shared by execution environment backends.
|
||||
|
||||
Kept separate from the base environment class so lazy backend imports do not
|
||||
depend on newly added exports from a large module cached earlier in a long-lived
|
||||
process.
|
||||
"""
|
||||
"""Path-component helpers shared by execution environment backends. Kept separate from the
|
||||
base class so lazy backend imports do not depend on newly added exports from a large module
|
||||
cached earlier in a long-lived process."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
@@ -20,23 +17,16 @@ _SANDBOX_DIR_HASH_LEN = 12
|
||||
|
||||
|
||||
def sanitize_task_id_for_path(task_id: str) -> str:
|
||||
"""Return a bind-mountable directory name for *task_id*'s sandbox.
|
||||
|
||||
Names that are already safe are returned verbatim, preserving existing
|
||||
sandbox locations. Rewritten names carry a digest because substitution
|
||||
alone is not injective: ``a:b`` and ``a_b`` must not share state.
|
||||
"""
|
||||
"""Bind-mountable directory name for *task_id*'s sandbox. Already-safe names are returned
|
||||
verbatim (preserving existing sandbox locations); rewritten names carry a digest because
|
||||
substitution alone is not injective: ``a:b`` and ``a_b`` must not share state."""
|
||||
value = task_id if isinstance(task_id, str) else ""
|
||||
if not value:
|
||||
return "default"
|
||||
|
||||
cleaned = _SANDBOX_DIR_UNSAFE_RE.sub("_", value)
|
||||
if (
|
||||
cleaned == value
|
||||
and len(value) <= _SANDBOX_DIR_MAX_LEN
|
||||
and value not in {".", ".."}
|
||||
and not value.endswith((".", " "))
|
||||
):
|
||||
if (cleaned == value and len(value) <= _SANDBOX_DIR_MAX_LEN
|
||||
and value not in {".", ".."} and not value.endswith((".", " "))):
|
||||
return value
|
||||
|
||||
digest = hashlib.sha256(value.encode("utf-8")).hexdigest()[:_SANDBOX_DIR_HASH_LEN]
|
||||
|
||||
@@ -7,16 +7,11 @@ import subprocess
|
||||
|
||||
def run_capture(cmd: list[str], *, timeout: float, check: bool = False, env: dict | None = None,
|
||||
) -> subprocess.CompletedProcess:
|
||||
"""``subprocess.run`` with the backend-standard capture settings.
|
||||
|
||||
Text mode with utf-8/replace decoding and stdin closed (DEVNULL) so a CLI
|
||||
that unexpectedly prompts cannot hang the agent.
|
||||
"""
|
||||
"""``subprocess.run`` with the backend-standard capture settings: text mode with utf-8/replace
|
||||
decoding and stdin closed (DEVNULL) so a CLI that unexpectedly prompts cannot hang the agent."""
|
||||
return subprocess.run(
|
||||
cmd,
|
||||
capture_output=True, text=True, encoding="utf-8", errors="replace",
|
||||
timeout=timeout, check=check, stdin=subprocess.DEVNULL, env=env,
|
||||
)
|
||||
cmd, capture_output=True, text=True, encoding="utf-8", errors="replace",
|
||||
timeout=timeout, check=check, stdin=subprocess.DEVNULL, env=env)
|
||||
|
||||
|
||||
def bash_argv(cmd_string: str, login: bool = False) -> list[str]:
|
||||
@@ -25,11 +20,9 @@ def bash_argv(cmd_string: str, login: bool = False) -> list[str]:
|
||||
|
||||
|
||||
def ensure_lazy_dep(feature: str) -> None:
|
||||
"""Lazy-install an optional SDK via ``tools.lazy_deps`` (idempotent).
|
||||
|
||||
Missing ``tools.lazy_deps`` is tolerated (the SDK import that follows will
|
||||
fail with its own message); any other failure surfaces as ``ImportError``.
|
||||
"""
|
||||
"""Lazy-install an optional SDK via ``tools.lazy_deps`` (idempotent). Missing ``tools.lazy_deps``
|
||||
is tolerated (the SDK import that follows fails with its own message); any other failure
|
||||
surfaces as ``ImportError``."""
|
||||
try:
|
||||
from tools.lazy_deps import ensure as _lazy_ensure
|
||||
_lazy_ensure(feature, prompt=False)
|
||||
|
||||
@@ -25,15 +25,13 @@ _SNAPSHOT_STORE = get_hermes_home() / "singularity_snapshots.json"
|
||||
|
||||
def _find_singularity_executable() -> str:
|
||||
"""Locate the apptainer or singularity CLI binary."""
|
||||
if shutil.which("apptainer"):
|
||||
return "apptainer"
|
||||
if shutil.which("singularity"):
|
||||
return "singularity"
|
||||
for exe in ("apptainer", "singularity"):
|
||||
if shutil.which(exe):
|
||||
return exe
|
||||
raise RuntimeError(
|
||||
"Neither 'apptainer' nor 'singularity' was found in PATH. "
|
||||
"Install Apptainer (https://apptainer.org/docs/admin/main/installation.html) "
|
||||
"or Singularity and ensure the CLI is available."
|
||||
)
|
||||
"or Singularity and ensure the CLI is available.")
|
||||
|
||||
|
||||
def _ensure_singularity_available() -> str:
|
||||
@@ -45,7 +43,6 @@ def _ensure_singularity_available() -> str:
|
||||
raise RuntimeError(f"Singularity backend selected but '{exe}' could not be executed.")
|
||||
except subprocess.TimeoutExpired:
|
||||
raise RuntimeError(f"'{exe} version' timed out.")
|
||||
|
||||
if result.returncode != 0:
|
||||
stderr = result.stderr.strip()[:200]
|
||||
raise RuntimeError(f"'{exe} version' failed (exit code {result.returncode}): {stderr}")
|
||||
@@ -61,24 +58,20 @@ def _save_snapshots(data: dict) -> None:
|
||||
|
||||
|
||||
def _get_scratch_dir() -> Path:
|
||||
"""``TERMINAL_SCRATCH_DIR`` override, else a writable ``/scratch`` (HPC), else the sandbox dir."""
|
||||
custom_scratch = os.getenv("TERMINAL_SCRATCH_DIR")
|
||||
if custom_scratch:
|
||||
scratch_path = Path(custom_scratch)
|
||||
scratch_path.mkdir(parents=True, exist_ok=True)
|
||||
return scratch_path
|
||||
|
||||
from tools.environments.base import get_sandbox_dir
|
||||
sandbox = get_sandbox_dir() / "singularity"
|
||||
|
||||
scratch = Path("/scratch")
|
||||
if scratch.exists() and os.access(scratch, os.W_OK):
|
||||
user_scratch = scratch / os.getenv("USER", "hermes") / "hermes-agent"
|
||||
user_scratch.mkdir(parents=True, exist_ok=True)
|
||||
logger.info("Using /scratch for sandboxes: %s", user_scratch)
|
||||
return user_scratch
|
||||
|
||||
sandbox.mkdir(parents=True, exist_ok=True)
|
||||
return sandbox
|
||||
else:
|
||||
from tools.environments.base import get_sandbox_dir
|
||||
scratch_path = get_sandbox_dir() / "singularity"
|
||||
scratch = Path("/scratch")
|
||||
if scratch.exists() and os.access(scratch, os.W_OK):
|
||||
scratch_path = scratch / os.getenv("USER", "hermes") / "hermes-agent"
|
||||
scratch_path.mkdir(parents=True, exist_ok=True)
|
||||
logger.info("Using /scratch for sandboxes: %s", scratch_path)
|
||||
scratch_path.mkdir(parents=True, exist_ok=True)
|
||||
return scratch_path
|
||||
|
||||
|
||||
def _get_apptainer_cache_dir() -> Path:
|
||||
@@ -92,9 +85,8 @@ _sif_build_lock = threading.Lock()
|
||||
|
||||
|
||||
def _get_or_build_sif(image: str, executable: str = "apptainer") -> str:
|
||||
if image.endswith('.sif') and Path(image).exists():
|
||||
return image
|
||||
if not image.startswith('docker://'):
|
||||
"""Build (once, cached) a SIF from a ``docker://`` URL; falls back to the URL on failure."""
|
||||
if (image.endswith('.sif') and Path(image).exists()) or not image.startswith('docker://'):
|
||||
return image
|
||||
|
||||
image_name = image.replace('docker://', '').replace('/', '-').replace(':', '-')
|
||||
@@ -160,11 +152,10 @@ class SingularityEnvironment(BaseEnvironment):
|
||||
self._memory = memory
|
||||
|
||||
if self._persistent:
|
||||
overlay_base = _get_scratch_dir() / "hermes-overlays"
|
||||
overlay_base.mkdir(parents=True, exist_ok=True)
|
||||
# A raw session-key task_id carries colons etc. unsafe in host path components;
|
||||
# the shared sanitizer keeps all backends agreeing on the mapping.
|
||||
self._overlay_dir = overlay_base / f"overlay-{sanitize_task_id_for_path(task_id)}"
|
||||
self._overlay_dir = (
|
||||
_get_scratch_dir() / "hermes-overlays" / f"overlay-{sanitize_task_id_for_path(task_id)}")
|
||||
self._overlay_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
self._start_instance()
|
||||
@@ -179,10 +170,8 @@ class SingularityEnvironment(BaseEnvironment):
|
||||
|
||||
try:
|
||||
from tools.credential_files import get_credential_file_mounts, get_skills_directory_mount
|
||||
for mount_entry in get_credential_file_mounts():
|
||||
cmd.extend(["--bind", f"{mount_entry['host_path']}:{mount_entry['container_path']}:ro"])
|
||||
for skills_mount in get_skills_directory_mount():
|
||||
cmd.extend(["--bind", f"{skills_mount['host_path']}:{skills_mount['container_path']}:ro"])
|
||||
for entry in (*get_credential_file_mounts(), *get_skills_directory_mount()):
|
||||
cmd.extend(["--bind", f"{entry['host_path']}:{entry['container_path']}:ro"])
|
||||
except Exception as e:
|
||||
logger.debug("Singularity: could not load credential/skills mounts: %s", e)
|
||||
|
||||
@@ -194,12 +183,12 @@ class SingularityEnvironment(BaseEnvironment):
|
||||
|
||||
try:
|
||||
result = run_capture(cmd, timeout=120)
|
||||
if result.returncode != 0:
|
||||
raise RuntimeError(f"Failed to start instance: {result.stderr}")
|
||||
self._instance_started = True
|
||||
logger.info("Singularity instance %s started (persistent=%s)", self.instance_id, self._persistent)
|
||||
except subprocess.TimeoutExpired:
|
||||
raise RuntimeError("Instance start timed out")
|
||||
if result.returncode != 0:
|
||||
raise RuntimeError(f"Failed to start instance: {result.stderr}")
|
||||
self._instance_started = True
|
||||
logger.info("Singularity instance %s started (persistent=%s)", self.instance_id, self._persistent)
|
||||
|
||||
def _run_bash(self, cmd_string: str, *, login: bool = False, timeout: int = 120,
|
||||
stdin_data: str | None = None) -> subprocess.Popen:
|
||||
|
||||
Reference in New Issue
Block a user