From 624a752e79f6057b46381aaedece8ed96f0ca081 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 10:36:07 -0700 Subject: [PATCH] =?UTF-8?q?refactor(agent/runtime):=20deadline/file=5Fsafe?= =?UTF-8?q?ty/estop/lifecycle=20leaf=20modules=20=E2=80=94=20dedupe=20resu?= =?UTF-8?q?lt=20plumbing,=20drop=20dead=20guard=20code?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - deadline: _result()/_abandon() replace 7 BoundedResult constructions and 2 cancel+callback sites; timers handled as a list; dead raise_if_timed_out removed. - file_safety: retired classify_cross_profile_target (0 refs; get_cross_profile_warning stub kept for external callers), _home_and_resolved/_mirror_warning shared by the sandbox/container mirror guards; _find_sandbox_mirror_segments inlined. - estop: _hermes_home/_canonical_root now the file_safety helpers; _reset_log_state_for_tests inlined into its only test. - subagent_lifecycle: _validate_request driven by _UNSUPPORTED_REQUEST_FIELDS table. - turn_liveness: dead start()/_abort_message removed, _emit_warning shared. - process_bootstrap: _enable_happy_eyeballs reused by the client variant. - Comment/docstring compaction across the remaining leaf modules. --- agent/__init__.py | 7 +- agent/battery.py | 53 +- agent/deadline.py | 440 +++-------- agent/delegation_context.py | 126 +-- agent/errors.py | 3 - agent/estop.py | 122 +-- agent/file_safety.py | 725 ++++++------------ agent/kanban_stop.py | 47 +- agent/oneshot.py | 52 +- agent/process_bootstrap.py | 206 ++--- agent/reactions.py | 28 +- agent/session_activity.py | 56 +- agent/subagent_lifecycle.py | 282 +++---- agent/thread_scoped_output.py | 59 +- agent/turn_liveness.py | 278 ++----- tests/agent/test_deadline.py | 5 +- tests/agent/test_file_safety_cross_profile.py | 57 +- .../agent/test_file_safety_sandbox_mirror.py | 4 - tests/test_estop.py | 4 +- 19 files changed, 706 insertions(+), 1848 deletions(-) diff --git a/agent/__init__.py b/agent/__init__.py index 41136f9b63..77f1cad846 100644 --- a/agent/__init__.py +++ b/agent/__init__.py @@ -1,8 +1,3 @@ -"""Agent internals -- extracted modules from run_agent.py. - -These modules contain pure utility functions and self-contained classes -that were previously embedded in the 3,600-line run_agent.py. Extracting -them makes run_agent.py focused on the AIAgent orchestrator class. -""" +"""Agent internals extracted from run_agent.py so it stays focused on AIAgent.""" from . import jiter_preload as _jiter_preload # noqa: F401 diff --git a/agent/battery.py b/agent/battery.py index a1c0f32fa4..9343c69d27 100644 --- a/agent/battery.py +++ b/agent/battery.py @@ -1,13 +1,9 @@ """System-battery read-out for the CLI/TUI status bar. -Reads the host battery through ``psutil`` (already a Hermes dependency) and -exposes a compact, colour-coded label. Everything degrades to "unavailable" -when there is no battery (desktops, servers, VMs) or when the read fails, so -callers can render the result unconditionally and simply show nothing. - -The status bar repaints often (every keystroke and on a ~1s idle refresh), so -:func:`read_battery` memoises the last reading for a few seconds instead of -hitting ``psutil`` on every frame. +Reads the host battery through ``psutil`` and exposes a compact, colour-coded +label. Everything degrades to "unavailable" (no battery / read failure) so +callers can render unconditionally. The status bar repaints on every keystroke, +so :func:`read_battery` memoises the reading for a few seconds. """ from __future__ import annotations @@ -19,12 +15,8 @@ from typing import Optional @dataclass(frozen=True) class BatteryStatus: - """A single battery reading. - - ``available`` is False on machines without a battery (or when the read - failed). ``percent`` is clamped to 0-100. ``plugged`` is True when on AC - power, False on battery, and None when the platform can't tell. - """ + """One reading: ``percent`` clamped 0-100; ``plugged`` None when the + platform can't tell.""" available: bool percent: Optional[int] = None @@ -45,6 +37,9 @@ CATEGORY_BAD = "bad" CATEGORY_CRITICAL = "critical" CATEGORY_DIM = "dim" +# (upper bound inclusive, category) for a discharging battery; first match wins. +_LEVEL_CATEGORIES = ((10, CATEGORY_CRITICAL), (20, CATEGORY_BAD), (50, CATEGORY_WARN)) + _CACHE_TTL_SECONDS = 8.0 _cache: Optional[tuple[float, BatteryStatus]] = None @@ -52,19 +47,11 @@ _cache: Optional[tuple[float, BatteryStatus]] = None def _read_battery_uncached() -> BatteryStatus: try: import psutil + + # ``sensors_battery`` is missing on some platforms/builds of psutil. + batt = getattr(psutil, "sensors_battery")() except Exception: return UNAVAILABLE - - # ``sensors_battery`` is missing on some platforms/builds of psutil. - reader = getattr(psutil, "sensors_battery", None) - if reader is None: - return UNAVAILABLE - - try: - batt = reader() - except Exception: - return UNAVAILABLE - if batt is None: return UNAVAILABLE @@ -75,12 +62,8 @@ def _read_battery_uncached() -> BatteryStatus: percent = max(0, min(100, int(round(float(raw_percent))))) except (TypeError, ValueError): percent = None - plugged = getattr(batt, "power_plugged", None) - if plugged is not None: - plugged = bool(plugged) - - return BatteryStatus(available=True, percent=percent, plugged=plugged) + return BatteryStatus(available=True, percent=percent, plugged=None if plugged is None else bool(plugged)) def read_battery(use_cache: bool = True) -> BatteryStatus: @@ -109,13 +92,9 @@ def battery_category(status: BatteryStatus) -> str: # On AC power the level isn't a concern — always read as healthy. if status.charging: return CATEGORY_GOOD - pct = status.percent - if pct <= 10: - return CATEGORY_CRITICAL - if pct <= 20: - return CATEGORY_BAD - if pct <= 50: - return CATEGORY_WARN + for bound, category in _LEVEL_CATEGORIES: + if status.percent <= bound: + return category return CATEGORY_GOOD diff --git a/agent/deadline.py b/agent/deadline.py index e3357f8a24..ecd9bfc306 100644 --- a/agent/deadline.py +++ b/agent/deadline.py @@ -1,63 +1,21 @@ """Unified deadline layer — one bounded-execution primitive, one timeout resolver. -Phase 1 of the architectural fix for the timeout/hang backlog -(https://github.com/NousResearch/hermes-agent/issues/85125). +Shared foundation for the site-local deadline mechanisms (#85125): -The tree currently carries at least six site-local deadline mechanisms, each -built for one incident, none shared (tool_executor batch deadline, telegram -``_await_with_thread_deadline``, gateway turn lease, reasoning stale floors, -``human_wait_ceiling``, per-MCP-handler timeouts). Every new stall report -grows that list by one. This module is the shared foundation the call sites -migrate onto in later phases: +* :func:`resolve_timeout` — config-first timeout resolution + (``timeouts:`` in config.yaml > legacy env var > default). +* :func:`clamp_timeout` — platform-safe clamping (huge timeouts overflow + ``time_t`` in ``Lock.acquire`` / ``Thread.join`` on macOS, #83220). +* :func:`run_bounded_async` — wall-clock deadline for awaitables driven by a + daemon ``threading.Timer``, so a blocked event loop cannot disable it (the + telegram adapter's ``_await_with_thread_deadline`` generalized). +* :func:`run_bounded_sync` — same contract for synchronous callables. +* :func:`kill_process_tree` — portable whole-tree termination. -* :func:`resolve_timeout` — one config-first resolution path for timeout - values (``timeouts:`` section in config.yaml > legacy env var > default), - so new surfaces stop inventing ``HERMES_*_TIMEOUT`` env vars (".env is for - secrets only") and hardcoded literals stop ignoring user config - (#63302, #53161, #43272 class). - -* :func:`clamp_timeout` — platform-safe clamping. Large user-supplied - timeouts overflow ``time_t`` inside ``threading.Lock.acquire(timeout=...)`` - / ``Thread.join(timeout=...)`` on macOS and kill whole tool batches - (#83220). Clamping at the shared boundary fixes that class once, for - every consumer. - -* :func:`run_bounded_async` — a wall-clock deadline for awaitables that does - NOT depend on event-loop timers. ``asyncio.wait_for`` schedules its expiry - on the loop; when the loop thread itself is blocked in a synchronous call - (family A of the #84047 stall triage), every asyncio-based timeout in the - process is silently disabled. This helper drives the deadline from a - daemon ``threading.Timer`` (generalizing the proven telegram-adapter - primitive) and abandons cancellation-shielded tasks instead of waiting for - cancellation to complete. The telegram adapter's private copy - (``plugins/platforms/telegram/adapter.py:_await_with_thread_deadline``) - migrates onto this in Phase 2 of #85125 — do not let the two drift in the - meantime; fix bugs here first. - -* :func:`run_bounded_sync` — the same contract for synchronous callables - bounded from a synchronous context (daemon worker thread, abandoned on - expiry). - -* :func:`kill_process_tree` — portable whole-tree termination so - kill-on-timeout stops orphaning descendants (#71148, #59549, #84967, - #68139 class). Existing site-local tree-kills that migrate onto this in - Phase 4 of #85125: ``gateway/status.py`` (taskkill wrapper + psutil - snapshot/reap pair) and ``tools/code_execution_tool.py`` (psutil - recursive children kill). - -Design invariants: - -* Exceptions raised by the bounded operation propagate unchanged — callers - keep their existing error handling. Only the *timeout* outcome is - reified (as :class:`BoundedResult`), because that is the outcome the - call sites keep getting wrong. -* A timeout produced by this layer is OUR deadline, not the provider's. - Callers that feed errors into ``agent/error_classifier.py`` should - classify :class:`DeadlineExpired` distinctly from transport timeouts - (the #59549 / #80323 misattribution class). -* ``None`` timeout means unbounded, and non-positive resolved values are - normalized to ``None`` (matching the existing - ``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` convention). +Invariants: operation exceptions propagate unchanged (only the *timeout* +outcome is reified as :class:`BoundedResult`); a timeout from this layer is +OUR deadline, not the provider's (classify :class:`DeadlineExpired` distinctly +from transport timeouts); ``None`` / non-positive timeout means unbounded. """ from __future__ import annotations @@ -87,35 +45,20 @@ __all__ = [ "kill_process_tree", ] -# Upper bound for any timeout handed to platform wait primitives. -# -# CPython converts ``threading.Lock.acquire(timeout=...)`` / -# ``Thread.join(timeout=...)`` deadlines to an absolute timestamp; very large -# relative timeouts overflow ``time_t`` on macOS and raise -# ``OverflowError: timestamp out of range for platform time_t`` (#83220). -# One year is semantically "unbounded" for every wait in this codebase while -# staying far below any platform conversion limit. -MAX_SAFE_TIMEOUT_S = 31_536_000.0 # 365 days +# One year: semantically "unbounded" yet far below any platform time_t limit +# (#83220: larger relative timeouts overflow inside Lock.acquire/Thread.join on macOS). +MAX_SAFE_TIMEOUT_S = 31_536_000.0 -# Grace period after a deadline fires before concluding the event loop thread -# is blocked in a synchronous call and dumping stacks (family A diagnostics). +# Grace after a deadline fires before concluding the loop thread is blocked and dumping stacks. _LOOP_BLOCKED_DUMP_GRACE_S = 5.0 -# ``Event.wait`` is a C-level block: KeyboardInterrupt / SetAsyncExc only -# land when the thread returns to Python. Slice the wait so a /stop or -# SIGINT during a bounded sync call is observed within this window rather -# than at the full deadline (#94285, tools/test_local_interrupt_cleanup). +# ``Event.wait`` is a C-level block: KeyboardInterrupt / SetAsyncExc only land when the +# thread returns to Python, so the sync wait is sliced to observe /stop or SIGINT promptly. _BOUNDED_SYNC_WAIT_SLICE_S = 0.2 class DeadlineExpired(TimeoutError): - """A deadline enforced by this layer expired. - - Distinct from transport/provider timeout types on purpose: when this is - raised (or a :class:`BoundedResult` reports ``timed_out``), the timeout - was Hermes's own bound — error classification must not attribute it to - the provider (#59549 / #80323 misattribution class). - """ + """A deadline enforced by this layer expired (Hermes's own bound, not the provider's).""" def __init__(self, label: str, timeout_s: float): super().__init__(f"deadline expired after {timeout_s:.1f}s: {label}") @@ -124,21 +67,13 @@ class DeadlineExpired(TimeoutError): class SuspectableBackend(Protocol): - """Phase 3a (#85125): a stateful backend the deadline layer can flag. + """A stateful backend (MCP connection, browser session, LSP client) the deadline layer can flag. - A timed-out stateful backend (MCP connection, browser session, LSP - client) may be left wedged by the abandoned half-finished operation. - ``run_bounded_*`` calls ``mark_suspect`` on timeout so the OWNER can - health-check or recycle the backend before reuse (``ensure_healthy``) - instead of returning a poisoned handle to the cache. Consumers adopt - incrementally (Phase 3b, one backend per PR), so the layer fails open: - backends without the protocol are simply never marked. - - Adopter contract: ``mark_suspect`` MUST be cheap, non-blocking, and - must not acquire locks the guarded operation may hold. It runs inline — - on the event loop in the async flavor, and on the caller's thread in - the sync flavor while the wedged worker is still alive. Set a flag; - do the expensive health-check/recycle work in ``ensure_healthy``. + ``run_bounded_*`` calls ``mark_suspect`` on timeout so the owner can health-check or + recycle the backend (``ensure_healthy``) before reuse. Backends without the protocol + are simply never marked. ``mark_suspect`` MUST be cheap, non-blocking, and must not + acquire locks the guarded operation may hold — it runs inline on the event loop / + caller's thread while the wedged worker is still alive. """ def mark_suspect(self, reason: str) -> None: ... @@ -147,13 +82,7 @@ class SuspectableBackend(Protocol): def _mark_backend_suspect(backend: object | None, label: str, timeout_s: float) -> None: - """Best-effort ``mark_suspect`` on a timed-out call's backend. - - Never raises: adoption state must not be able to weaken the deadline - bound or corrupt the ``BoundedResult`` the caller is about to receive. - A non-adopting backend (no ``mark_suspect``) is tolerated silently — - Phase 3b lands per-backend, so absence is the norm during adoption. - """ + """Best-effort ``mark_suspect``; never raises, non-adopting backends tolerated.""" if backend is None: return try: @@ -166,12 +95,7 @@ def _mark_backend_suspect(backend: object | None, label: str, timeout_s: float) @dataclass(frozen=True, kw_only=True) class BoundedResult: - """Outcome of a bounded operation. - - ``timed_out`` is the reified outcome; on completion ``value`` holds the - operation's return value. Operation exceptions are never captured here — - they propagate to the caller unchanged. - """ + """Outcome of a bounded operation; operation exceptions are never captured here.""" timed_out: bool value: Any @@ -179,33 +103,21 @@ class BoundedResult: timeout_s: Optional[float] label: str - def raise_if_timed_out(self) -> Any: - """Return ``value``, raising :class:`DeadlineExpired` on timeout.""" - if self.timed_out: - raise DeadlineExpired(self.label, float(self.timeout_s or 0.0)) - return self.value + +def _result(start: float, timeout_s: Optional[float], label: str, *, value: Any = None, timed_out: bool = False) -> BoundedResult: + return BoundedResult( + timed_out=timed_out, value=value, elapsed_s=time.monotonic() - start, timeout_s=timeout_s, label=label + ) def clamp_timeout(timeout: Optional[float]) -> Optional[float]: - """Normalize a timeout value for platform wait primitives. - - * ``None`` stays ``None`` (unbounded). - * Non-positive values become ``None`` (unbounded) — matching the existing - ``HERMES_CONCURRENT_TOOL_TIMEOUT_S`` "0 disables the bound" convention. - * Values above :data:`MAX_SAFE_TIMEOUT_S` are capped so they can never - overflow ``time_t`` inside ``Lock.acquire`` / ``Thread.join`` on macOS - (#83220). - * Non-numeric values are treated as unset (``None``) with a warning - rather than crashing the call path they were meant to protect. - """ + """Normalize a timeout: None/non-positive/non-numeric/NaN -> None (unbounded), else capped.""" if timeout is None: return None try: value = float(timeout) except (TypeError, ValueError): - logger.warning( - "clamp_timeout: non-numeric timeout %r; treating as unbounded", timeout - ) + logger.warning("clamp_timeout: non-numeric timeout %r; treating as unbounded", timeout) return None if value != value: # NaN logger.warning("clamp_timeout: NaN timeout; treating as unbounded") @@ -215,18 +127,11 @@ def clamp_timeout(timeout: Optional[float]) -> Optional[float]: return min(value, MAX_SAFE_TIMEOUT_S) -# --------------------------------------------------------------------------- -# Timeout resolution: config.yaml ``timeouts:`` section > legacy env var > -# registered default. -# --------------------------------------------------------------------------- +# --- Timeout resolution: config ``timeouts:`` > legacy env var > default ------ def _timeouts_section() -> dict: - """Read the ``timeouts:`` root section from config.yaml (read-only). - - Isolated for testability and so a broken config read can never take down - the call path the timeout was protecting. - """ + """Read the ``timeouts:`` root section from config.yaml (read-only, fail-open).""" try: from hermes_cli.config import load_config_readonly @@ -253,30 +158,16 @@ def resolve_timeout( default: Optional[float], env_var: Optional[str] = None, ) -> Optional[float]: - """Resolve a timeout in seconds for a dotted config key. + """Resolve a timeout (seconds) for dotted ``timeouts.`` > ``env_var`` > ``default``. - Precedence (established by the ``providers.*.request_timeout_seconds`` - pattern — config wins over the legacy env var): - - 1. ``timeouts.`` in config.yaml (dotted key walks nested maps, e.g. - ``tools.concurrent_batch`` reads ``timeouts: {tools: {concurrent_batch: ...}}``) - 2. ``env_var`` when set and non-empty (legacy bridge — internal mechanism - and back-compat only; new surfaces must not grow new user-facing - ``HERMES_*`` timeout env vars) - 3. ``default`` - - The winning value is passed through :func:`clamp_timeout`, so ``0`` or a - negative value means "unbounded" and oversized values are made - platform-safe. Invalid (non-numeric) config/env values fall through to - the next source with a warning instead of breaking the protected path. + The winner goes through :func:`clamp_timeout`; invalid config/env values fall + through to the next source with a warning. """ raw = _lookup_dotted(_timeouts_section(), key) if raw is not None: - # Explicit float() (clamp_timeout would also convert) so that invalid - # config values FALL THROUGH to the env var / default instead of - # resolving as unbounded — do not "simplify" this away. bool is - # rejected because YAML `true` would silently become a 1-second - # deadline; NaN is rejected for the same fall-through reason. + # Explicit float() so invalid config values FALL THROUGH to env/default instead of + # resolving as unbounded. bool rejected (YAML `true` would become a 1s deadline); + # NaN rejected for the same fall-through reason. if not isinstance(raw, bool): try: value = float(raw) @@ -284,9 +175,7 @@ def resolve_timeout( return clamp_timeout(value) except (TypeError, ValueError): pass - logger.warning( - "timeouts.%s: invalid value %r in config.yaml; ignoring", key, raw - ) + logger.warning("timeouts.%s: invalid value %r in config.yaml; ignoring", key, raw) if env_var: env_raw = os.getenv(env_var, "").strip() @@ -299,15 +188,9 @@ def resolve_timeout( return clamp_timeout(default) -# --------------------------------------------------------------------------- -# Bounded execution — async flavor. -# -# Generalizes plugins/platforms/telegram/adapter.py:_await_with_thread_deadline -# (the #63309 fix): the deadline is driven by a daemon threading.Timer so a -# blocked event loop cannot disable it, and a second timer dumps all thread -# stacks when the loop provably failed to process the expiry — the one piece -# of information loop-blocked hangs otherwise never surface. -# --------------------------------------------------------------------------- +# --- Bounded execution — async flavor ------------------------------------------ +# The deadline is a daemon threading.Timer so a blocked event loop cannot disable it; a +# second timer dumps all thread stacks when the loop provably failed to process the expiry. def _consume_abandoned(task: "asyncio.Future[Any]") -> None: @@ -319,8 +202,14 @@ def _consume_abandoned(task: "asyncio.Future[Any]") -> None: pass +def _abandon(task: "asyncio.Future[Any]") -> None: + """Cancel ``task`` and never await it; its outcome is consumed so it stays unobserved-safe.""" + task.cancel() + task.add_done_callback(_consume_abandoned) + + async def _run_abandon_cleanup(on_abandon: Callable[[], Awaitable[Any]]) -> None: - """Run abandonment cleanup fully fire-and-forget (its failures swallowed).""" + """Run abandonment cleanup fire-and-forget (its failures swallowed).""" try: await on_abandon() except Exception: @@ -355,30 +244,15 @@ async def run_bounded_async( ) -> BoundedResult: """Await ``awaitable`` under a wall-clock deadline independent of loop timers. - On completion returns ``BoundedResult(timed_out=False, value=...)``; - exceptions from the operation (including ``asyncio.CancelledError`` from a - caller cancelling *us*) propagate unchanged. - - On timeout the underlying task is cancelled and **abandoned** — we do not - await cancellation completion, because cancellation-shielded scopes (anyio, - httpcore init, MCP SDK teardown) are exactly the paths that wedge forever. - ``on_abandon`` (zero-arg callable returning an awaitable) is scheduled as - detached best-effort cleanup for the half-built state the abandoned task - may leave behind. Returns ``BoundedResult(timed_out=True, value=None)``. - - ``timeout=None`` (or a non-positive resolved value) awaits unbounded. + Operation exceptions (incl. ``CancelledError`` from a caller cancelling *us*) + propagate unchanged. On timeout the task is cancelled and **abandoned** (never + awaited — cancellation-shielded scopes are exactly the paths that wedge), and + ``on_abandon`` is scheduled as detached best-effort cleanup. """ timeout_s = clamp_timeout(timeout) start = time.monotonic() if timeout_s is None: - value = await awaitable - return BoundedResult( - timed_out=False, - value=value, - elapsed_s=time.monotonic() - start, - timeout_s=None, - label=label, - ) + return _result(start, None, label, value=await awaitable) task = asyncio.ensure_future(awaitable) loop = asyncio.get_running_loop() @@ -390,84 +264,44 @@ async def run_bounded_async( if not deadline.done(): deadline.set_result(None) - def _expire_from_thread() -> None: - loop.call_soon_threadsafe(_mark_expired) - def _watchdog_check() -> None: if not loop_processed_expiry.is_set(): _dump_blocked_loop_diagnostics(label, timeout_s) - timer = threading.Timer(timeout_s, _expire_from_thread) - timer.daemon = True - timer.start() - watchdog: Optional[threading.Timer] = None + timers = [threading.Timer(timeout_s, lambda: loop.call_soon_threadsafe(_mark_expired))] if dump_on_blocked_loop: - watchdog = threading.Timer( - timeout_s + _LOOP_BLOCKED_DUMP_GRACE_S, _watchdog_check - ) - watchdog.daemon = True - watchdog.start() + timers.append(threading.Timer(timeout_s + _LOOP_BLOCKED_DUMP_GRACE_S, _watchdog_check)) + for t in timers: + t.daemon = True + t.start() try: try: - done, _ = await asyncio.wait( - {task, deadline}, return_when=asyncio.FIRST_COMPLETED - ) + done, _ = await asyncio.wait({task, deadline}, return_when=asyncio.FIRST_COMPLETED) except asyncio.CancelledError: - # The CALLER cancelled us. Without this, `task` would keep running - # unobserved (and later log "exception was never retrieved") — - # a leak the telegram original also had. Cancel + abandon it, then - # let the cancellation propagate. - task.cancel() - task.add_done_callback(_consume_abandoned) + _abandon(task) # the CALLER cancelled us; `task` must not run unobserved raise if task in done: if not deadline.done(): deadline.cancel() - value = await task - return BoundedResult( - timed_out=False, - value=value, - elapsed_s=time.monotonic() - start, - timeout_s=timeout_s, - label=label, - ) + return _result(start, timeout_s, label, value=await task) - task.cancel() - task.add_done_callback(_consume_abandoned) + _abandon(task) if on_abandon is not None: - cleanup = asyncio.ensure_future(_run_abandon_cleanup(on_abandon)) - cleanup.add_done_callback(_consume_abandoned) - # Phase 3a (#85125): the abandoned task may leave the backend - # half-wedged; flag it so the owner recycles before reuse. - # Deliberately INLINE on the loop (adopter contract: mark_suspect is - # cheap and non-blocking). Running it synchronously guarantees the - # mark happens-before this BoundedResult returns AND before the - # ensure_future'd on_abandon cleanup can start (next loop tick) — an - # offloaded mark would race both. + asyncio.ensure_future(_run_abandon_cleanup(on_abandon)).add_done_callback(_consume_abandoned) + # Deliberately INLINE on the loop: the mark must happen-before this result returns + # AND before the ensure_future'd on_abandon cleanup starts (next tick). _mark_backend_suspect(backend, label, timeout_s) - logger.warning( - "[deadline] %r timed out after %.1fs; task abandoned", label, timeout_s - ) - return BoundedResult( - timed_out=True, - value=None, - elapsed_s=time.monotonic() - start, - timeout_s=timeout_s, - label=label, - ) + logger.warning("[deadline] %r timed out after %.1fs; task abandoned", label, timeout_s) + return _result(start, timeout_s, label, timed_out=True) finally: - timer.cancel() - if watchdog is not None: - watchdog.cancel() - # cancel() cannot stop a Timer whose callback is already running; - # setting the event closes that race so a completed await can never - # be misreported as a blocked loop. + for t in timers: + t.cancel() + # cancel() cannot stop a Timer whose callback is already running; setting the + # event closes that race so a completed await is never misreported as blocked. loop_processed_expiry.set() -# --------------------------------------------------------------------------- -# Bounded execution — sync flavor. -# --------------------------------------------------------------------------- +# --- Bounded execution — sync flavor ------------------------------------------- def run_bounded_sync( @@ -480,37 +314,15 @@ def run_bounded_sync( ) -> BoundedResult: """Run ``fn`` in a daemon worker thread under a wall-clock deadline. - On completion returns its value (exceptions re-raised in the caller). - On expiry the worker thread is **abandoned** (daemon, so it cannot block - interpreter exit), ``on_timeout`` (if given) runs best-effort in the - caller's thread — e.g. to mark a backend suspect or kill a subprocess — - and ``BoundedResult(timed_out=True)`` is returned. - - Intended for infrequent, seconds-scale blocking backend calls. Do NOT - use per-item in hot loops: each call spawns a thread, and every timeout - permanently leaks an abandoned daemon thread — a wedged backend called - in a retry loop would accumulate them. - - ``timeout=None`` (or non-positive) blocks until ``fn`` returns. - - The worker runs under ``contextvars.copy_context()`` so profile secret - scope, session id, and delegated-child guards set on the caller survive - the thread hop (terminal env.execute — #94285 CI). - - The caller's wait is sliced (``_BOUNDED_SYNC_WAIT_SLICE_S``) so a - ``KeyboardInterrupt`` or ``PyThreadState_SetAsyncExc`` lands within - that window instead of only at the full deadline. + Exceptions re-raise in the caller. On expiry the worker is **abandoned** and + ``on_timeout`` runs best-effort in the caller's thread. Every timeout leaks one + daemon thread, so do NOT use per-item in hot loops. The worker runs under + ``contextvars.copy_context()`` so secret scope / session id survive the hop. """ timeout_s = clamp_timeout(timeout) start = time.monotonic() if timeout_s is None: - return BoundedResult( - timed_out=False, - value=fn(), - elapsed_s=time.monotonic() - start, - timeout_s=None, - label=label, - ) + return _result(start, None, label, value=fn()) box: dict[str, Any] = {} done = threading.Event() @@ -534,67 +346,35 @@ def run_bounded_sync( done.wait(min(_BOUNDED_SYNC_WAIT_SLICE_S, remaining)) if not done.is_set(): - logger.warning( - "[deadline] %r timed out after %.1fs; worker abandoned", label, timeout_s - ) - # Phase 3a (#85125), ordering: mark suspect BEFORE owner cleanup so a - # recycle/re-init in on_timeout never gets a stale flag on the healed - # replacement. The sync flavor runs the mark inline — the protocol - # contract requires mark_suspect to be cheap. + logger.warning("[deadline] %r timed out after %.1fs; worker abandoned", label, timeout_s) + # Mark suspect BEFORE owner cleanup so a recycle in on_timeout never + # inherits a stale flag on the healed replacement. _mark_backend_suspect(backend, label, timeout_s) if on_timeout is not None: try: on_timeout() except Exception: logger.debug("deadline on_timeout callback failed", exc_info=True) - return BoundedResult( - timed_out=True, - value=None, - elapsed_s=time.monotonic() - start, - timeout_s=timeout_s, - label=label, - ) + return _result(start, timeout_s, label, timed_out=True) if "exc" in box: raise box["exc"] - return BoundedResult( - timed_out=False, - value=box.get("value"), - elapsed_s=time.monotonic() - start, - timeout_s=timeout_s, - label=label, - ) + return _result(start, timeout_s, label, value=box.get("value")) -# --------------------------------------------------------------------------- -# Whole-tree process termination. -# --------------------------------------------------------------------------- +# --- Whole-tree process termination -------------------------------------------- def kill_process_tree(pid: int, *, sig: Optional[int] = None) -> bool: """Terminate ``pid`` and all its descendants, portably. - Kill-on-timeout that signals only the direct child orphans process trees - (cron scripts, in-container shells, browser daemons — #71148 class). + Windows: ``taskkill /F /T`` (``sig`` ignored). POSIX: descendants are snapshotted + via psutil BEFORE signalling (once the parent dies they reparent and a parent walk + finds nothing), then the process group is signalled when ``pid`` leads one, and + every snapshotted descendant individually — which also reaches children that + ``setsid`` into their own session. ``sig`` defaults to ``SIGKILL``. - * Windows: ``taskkill /F /T`` terminates the tree (``sig`` ignored; - Windows has no equivalent). Console-window flash is suppressed via - ``windows_hide_flags`` and the exit code is checked, so a dead or - inaccessible PID reports ``False`` like the POSIX path. - * POSIX: the descendant set is snapshotted via psutil (a hard - dependency) BEFORE any signal — once the parent dies its children are - reparented and can no longer be found by a parent walk. Then the - process group is signalled when ``pid`` leads one (covers - grandchildren in the same session in one syscall), and every - snapshotted descendant is signalled individually — which also reaches - descendants that created their OWN sessions (a child that called - ``setsid``, exactly what user shell commands do; see - tools/environments/base.py). ``sig`` defaults to ``SIGKILL``. - psutil's identity-aware ``Process`` (PID + create time) means a - recycled PID is never signalled. - - Returns True when the target (or any of its tree) was signalled, False - when the process was already gone or every termination call failed. + Returns True when the target (or any of its tree) was signalled. """ if sys.platform == "win32": try: @@ -611,13 +391,10 @@ def kill_process_tree(pid: int, *, sig: Optional[int] = None) -> bool: check=False, creationflags=creationflags, ) - # taskkill exits non-zero for not-found / access-denied; keep the - # cross-platform contract (False = nothing was terminated). + # taskkill exits non-zero for not-found / access-denied (False = nothing terminated). return proc.returncode == 0 except Exception: - logger.debug( - "kill_process_tree: taskkill failed for pid %s", pid, exc_info=True - ) + logger.debug("kill_process_tree: taskkill failed for pid %s", pid, exc_info=True) return False import signal as _signal @@ -625,34 +402,25 @@ def kill_process_tree(pid: int, *, sig: Optional[int] = None) -> bool: if sig is None: sig = _signal.SIGKILL - # Snapshot descendants while the parent is still alive — after it dies - # they reparent to init/subreaper and a parent walk finds nothing. - descendants: list = [] try: import psutil descendants = psutil.Process(int(pid)).children(recursive=True) except Exception: - # Already gone, or psutil unavailable in a stripped env — the - # group-signal below still covers same-session descendants. + # Already gone, or psutil unavailable — the group signal still covers same-session descendants. descendants = [] signalled = False try: - # NOTE: getpgid→killpg has an inherent TOCTOU (pid could be reaped and - # recycled between the calls). All existing killpg sites share it; the - # psutil sweep below is identity-aware and does not. + # getpgid→killpg has an inherent TOCTOU shared by every killpg site; the psutil + # sweep below is identity-aware (PID + create time) and does not. pgid = os.getpgid(pid) except (ProcessLookupError, PermissionError, OSError): pgid = None try: if pgid is not None and pgid == pid: - # pid leads its own group: one syscall covers the whole group. - # (The == check guards against signalling the caller's own group - # when pid is not a leader.) - os.killpg( # windows-footgun: ok — POSIX-only branch (win32 returns above) - pgid, sig - ) + # pid leads its own group (the == check avoids signalling the caller's group). + os.killpg(pgid, sig) # windows-footgun: ok — POSIX-only branch (win32 returns above) else: os.kill(pid, sig) signalled = True @@ -661,8 +429,6 @@ def kill_process_tree(pid: int, *, sig: Optional[int] = None) -> bool: except (PermissionError, OSError): logger.debug("kill_process_tree: signal failed for pid %s", pid, exc_info=True) - # Sweep the snapshot: reaches descendants outside the parent's group - # (their own setsid sessions) and the non-group-leader case. for child in descendants: try: if child.is_running(): # identity-aware: recycled PIDs skipped diff --git a/agent/delegation_context.py b/agent/delegation_context.py index 9b8bce9759..5ff3eb8e8c 100644 --- a/agent/delegation_context.py +++ b/agent/delegation_context.py @@ -1,36 +1,23 @@ """Context-local state for delegate_task child execution. -The parent Hermes process may itself be a Kanban dispatcher worker with -HERMES_KANBAN_* variables in process env. delegate_task children run inside the -same Python process, but they are not dispatcher-owned Kanban workers. This -module lets code paths that resolve tool schemas or spawn subprocesses fail -closed for delegated children without mutating global os.environ for the parent. - -Cron jobs need the same treatment for the same reason: ``cronjob(action="run")`` -executes ``run_job()`` in-process, so a cron agent fired from inside a Kanban -worker would otherwise inherit that worker's dispatcher identity. -``non_dispatcher_owned_context()`` covers both cases. +A Hermes process may itself be a Kanban dispatcher worker with HERMES_KANBAN_* +in os.environ. delegate_task children (in-process) and cron jobs fired via +``cronjob(action="run")`` are NOT dispatcher-owned, so identity gates must +fail closed for them without mutating the process-global environment. """ from __future__ import annotations +import os from contextlib import contextmanager from contextvars import ContextVar, Token from typing import Iterator, Mapping, MutableMapping -_DELEGATED_CHILD_CONTEXT: ContextVar[bool] = ContextVar( - "hermes_delegated_child_context", - default=False, -) +_DELEGATED_CHILD_CONTEXT: ContextVar[bool] = ContextVar("hermes_delegated_child_context", default=False) -# Set for any in-process execution that is NOT the dispatcher-owned worker even -# though the worker's HERMES_KANBAN_* vars are legitimately in os.environ (cron -# jobs fired via the `cronjob` tool). Kept separate from -# _DELEGATED_CHILD_CONTEXT so the delegate_task-specific behaviour attached to -# that flag (subprocess env scrubbing, its own error strings) is unchanged. -_NON_DISPATCHER_OWNED_CONTEXT: ContextVar[bool] = ContextVar( - "hermes_non_dispatcher_owned_context", - default=False, -) +# Any in-process execution that is NOT the dispatcher-owned worker (cron jobs). +# Kept separate from _DELEGATED_CHILD_CONTEXT so delegate_task-specific +# behaviour (subprocess env scrubbing, its error strings) is unchanged. +_NON_DISPATCHER_OWNED_CONTEXT: ContextVar[bool] = ContextVar("hermes_non_dispatcher_owned_context", default=False) DELEGATED_CHILD_ENV_MARKER = "HERMES_DELEGATED_CHILD_CONTEXT" @@ -49,14 +36,12 @@ KANBAN_ENV_KEYS: tuple[str, ...] = ( def delegated_child_context(session_id: str | None = None) -> Iterator[None]: """Mark child execution and isolate its task-local session identity. - Child construction calls ``set_current_session_id`` internally, so even a - context entered without an id must restore the parent's ContextVar. Child - execution passes its explicit id and receives it only for this scope. + Even a context entered without an id must restore the parent's session + ContextVar (child construction calls ``set_current_session_id``). """ token = _DELEGATED_CHILD_CONTEXT.set(True) try: - # Import lazily: session_context calls is_delegated_child_context() when - # deciding whether the compatibility os.environ mirror is safe. + # Lazy: session_context calls is_delegated_child_context(). from gateway.session_context import scoped_current_session_id with scoped_current_session_id(session_id): @@ -70,49 +55,8 @@ def is_delegated_child_context() -> bool: return bool(_DELEGATED_CHILD_CONTEXT.get()) -@contextmanager -def non_dispatcher_owned_context() -> Iterator[None]: - """Mark in-process execution that does NOT own the dispatcher's Kanban task. - - A Kanban worker is a normal CLI agent whose default toolset includes - ``cronjob``; ``cronjob(action="run")`` runs ``run_job()`` inside the worker's - own process, where ``HERMES_KANBAN_TASK`` is legitimately set. Without this - marker the cron agent is misread as that worker: the kanban toolset is - force-added, the worker protocol is injected into its system prompt, and - ``kanban_complete`` defaults ``task_id`` to ``$HERMES_KANBAN_TASK`` — letting - an unrelated cron job close the worker's task and overwrite real results. - - Scoped via ContextVar rather than by clearing ``os.environ``: the env is - process-global and shared with the worker's own claim heartbeat, the - gateway's Kanban watchers, and concurrent cron jobs on the parallel pool, so - mutating it would starve the worker's claim and race those readers. - """ - token = _NON_DISPATCHER_OWNED_CONTEXT.set(True) - try: - yield - finally: - _NON_DISPATCHER_OWNED_CONTEXT.reset(token) - - -def is_dispatcher_owned_worker_context() -> bool: - """Return True only when this execution owns the dispatcher's Kanban task. - - The single predicate every ``HERMES_KANBAN_*`` identity gate should use - before trusting those vars. False for delegate_task children and for cron - jobs fired in-process from a worker. - """ - if _DELEGATED_CHILD_CONTEXT.get(): - return False - return not _NON_DISPATCHER_OWNED_CONTEXT.get() - - def enter_non_dispatcher_owned_context() -> Token[bool]: - """Token-based form of :func:`non_dispatcher_owned_context`. - - For callers whose scope is a long ``try`` with a matching ``finally`` rather - than a ``with`` block (``cron.scheduler.run_job``). Pair with - :func:`exit_non_dispatcher_owned_context`. - """ + """Token form of :func:`non_dispatcher_owned_context` for long try/finally scopes.""" return _NON_DISPATCHER_OWNED_CONTEXT.set(True) @@ -121,13 +65,30 @@ def exit_non_dispatcher_owned_context(token: Token[bool]) -> None: _NON_DISPATCHER_OWNED_CONTEXT.reset(token) +@contextmanager +def non_dispatcher_owned_context() -> Iterator[None]: + """Mark in-process execution that does NOT own the dispatcher's Kanban task. + + Without it a cron agent run inside a worker is misread as that worker + (kanban toolset force-added, ``kanban_complete`` defaulting to its task). + ContextVar-scoped rather than clearing os.environ, which the worker's claim + heartbeat and concurrent readers share. + """ + token = enter_non_dispatcher_owned_context() + try: + yield + finally: + exit_non_dispatcher_owned_context(token) + + +def is_dispatcher_owned_worker_context() -> bool: + """The single predicate every ``HERMES_KANBAN_*`` identity gate should use.""" + return not (_DELEGATED_CHILD_CONTEXT.get() or _NON_DISPATCHER_OWNED_CONTEXT.get()) + + def is_delegated_child_process_context() -> bool: """Return True in this process or a subprocess spawned by a child.""" - import os - - return bool(_DELEGATED_CHILD_CONTEXT.get()) or bool( - os.environ.get(DELEGATED_CHILD_ENV_MARKER) - ) + return bool(_DELEGATED_CHILD_CONTEXT.get()) or bool(os.environ.get(DELEGATED_CHILD_ENV_MARKER)) def scrub_kanban_env(env: Mapping[str, str] | MutableMapping[str, str]) -> dict[str, str]: @@ -144,18 +105,9 @@ def delegated_child_subprocess_env( ) -> dict[str, str] | None: """Return an env override only when delegated-child lineage must cross fork. - Most subprocess call sites historically used ``env=None`` to inherit the - process environment. In a ``delegate_task`` child, inheriting as-is leaks - parent dispatcher ``HERMES_KANBAN_*`` vars while losing the ContextVar in - the new process. This helper preserves normal ``env=None`` semantics for - non-delegated calls, and only materializes a scrubbed env when the lineage - marker must be propagated across a child-process boundary. + Preserves ``env=None`` inherit semantics for non-delegated calls; in a + child, materializes a scrubbed env carrying the lineage marker. """ if not is_delegated_child_process_context(): return None if env is None else dict(env) - - if env is None: - import os - - env = os.environ - return scrub_kanban_env(env) + return scrub_kanban_env(os.environ if env is None else env) diff --git a/agent/errors.py b/agent/errors.py index 8b3db5a358..66d1c1329a 100644 --- a/agent/errors.py +++ b/agent/errors.py @@ -1,13 +1,10 @@ class SSLConfigurationError(Exception): """Raised when SSL/TLS certificate bundle configuration fails.""" - pass class EmptyStreamError(RuntimeError): """Raised when a provider closes a stream without yielding a response.""" - pass - class MoAPresetNotFoundError(ValueError): """Raised when a persisted MoA preset no longer exists in config.""" diff --git a/agent/estop.py b/agent/estop.py index b62fc378d6..59cb0ffbc4 100644 --- a/agent/estop.py +++ b/agent/estop.py @@ -1,78 +1,43 @@ """Global emergency stop (ESTOP) — a resumable pause for NEW work only. -``hermes pause`` writes a sentinel file at ``$HERMES_HOME/ESTOP``; -``hermes resume`` removes it. While the sentinel exists: - -* the cron scheduler skips dispatching due jobs (``cron/scheduler.py:tick``), -* the embedded kanban dispatcher skips spawning workers - (``gateway/kanban_watchers.py``), -* new gateway turns get a brief "Hermes is paused" reply instead of an - agent run (``gateway/run.py:_handle_message``). - -In-flight work is NEVER killed — this is pause-new-work, not panic/exit. -The check is one or two ``os.stat`` calls (process home + fleet root when -they differ) so callers may run it every tick; no caching beyond the OS is -performed, so engaging/disengaging takes effect on the very next check. - -The sentinel body is optional JSON ``{"reason": ..., "engaged_at": ...}``. -A corrupt or empty file still counts as engaged (fail safe): the pause must -hold even if the file was created by ``touch ~/.hermes/ESTOP``. - -Ported from: gastownhall/gastown estop.go (MIT). Related prior art: -#26778 (/panic — kill/exit semantics; deliberately different, ours is -resumable) and #44617 (interrupting in-flight cron; deliberately out of -scope here). +``hermes pause`` writes a sentinel at ``$HERMES_HOME/ESTOP``; ``hermes resume`` +removes it. While it exists the cron scheduler, kanban dispatcher and new +gateway turns skip work; in-flight work is never killed. The check is one or +two ``os.stat`` calls (process home + fleet root when they differ), uncached, +so engaging/disengaging binds on the next check. The sentinel body is optional +JSON ``{"reason", "engaged_at"}``; a corrupt/empty file still counts as engaged +(fail safe, e.g. ``touch ~/.hermes/ESTOP``). Ported from gastownhall/gastown +estop.go (MIT). """ from __future__ import annotations import json import logging -import os import threading from datetime import datetime, timezone from pathlib import Path from typing import Optional +# Same profile-aware / fleet-root resolvers the file-safety guards use (fail-open to ~/.hermes). +from agent.file_safety import _hermes_home_path as _hermes_home, _hermes_root_path as _canonical_root + SENTINEL_NAME = "ESTOP" -# Per-component "logged already for this engagement" flags so a paused -# dispatch loop logs once per engagement instead of once per tick. +# Per-component "logged already for this engagement" flags: log once per +# engagement, not once per tick. _log_lock = threading.Lock() _logged_components: set[str] = set() -def _hermes_home() -> Path: - """Resolve the active HERMES_HOME (profile-aware) at call time.""" - try: - from hermes_constants import get_hermes_home - return get_hermes_home() - except Exception: - return Path(os.path.expanduser("~/.hermes")) - - -def _canonical_root() -> Path: - """Fleet-wide Hermes root, even when this process is a profile gateway. - - Profile gateways launch with HERMES_HOME=~/.hermes/profiles/. - ``hermes pause`` from an operator seat writes ~/.hermes/ESTOP. If we - only inspect the profile home, the emergency stop does not bind - (jarvis-os/t_7b65ff88: fleet-analyst kept dispatching through pause). - """ - try: - from hermes_constants import get_default_hermes_root - return Path(get_default_hermes_root()) - except Exception: - return Path(os.path.expanduser("~/.hermes")) - - def sentinel_path() -> Path: """Path of the ESTOP sentinel this process would write on `hermes pause`.""" return _hermes_home() / SENTINEL_NAME def _candidate_sentinel_paths() -> list: - """Profile home first, then the fleet root if it is a different directory.""" + """Profile home first, then the fleet root if it is a different directory: a profile + gateway (HERMES_HOME=~/.hermes/profiles/) must still honor an operator's ~/.hermes/ESTOP.""" primary = sentinel_path() paths = [primary] try: @@ -83,20 +48,14 @@ def _candidate_sentinel_paths() -> list: if root.resolve() != primary.resolve(): paths.append(root) except Exception: - # Non-Path test doubles (fail-safe stat fixture) fail .resolve(); - # the generic comparison below still dedupes plain equal paths. + # Non-Path test doubles fail .resolve(); plain equality still dedupes. if root != primary: paths.append(root) return paths def is_engaged() -> bool: - """Cheap check: is the global emergency stop engaged? - - Engaged if ANY candidate sentinel exists: the process HERMES_HOME - (profile-local) or the fleet canonical root (~/.hermes). Fail SAFE on - stat errors so an unreadable sentinel still holds the pause. - """ + """True if ANY candidate sentinel exists; fail SAFE (True) on stat errors.""" saw_stat_error = False for path in _candidate_sentinel_paths(): try: @@ -110,10 +69,7 @@ def is_engaged() -> bool: def engage(reason: Optional[str] = None) -> Path: """Create the ESTOP sentinel. Idempotent; re-engaging updates the file.""" path = sentinel_path() - payload = { - "engaged_at": datetime.now(timezone.utc).isoformat(), - "reason": reason or None, - } + payload = {"engaged_at": datetime.now(timezone.utc).isoformat(), "reason": reason or None} try: path.parent.mkdir(parents=True, exist_ok=True) path.write_text(json.dumps(payload, indent=2) + "\n", encoding="utf-8") @@ -127,34 +83,25 @@ def engage(reason: Optional[str] = None) -> Path: def disengage() -> bool: - """Remove ESTOP sentinels this process can see. - - Lifts both the process-local sentinel and the fleet-root sentinel so - ``hermes resume`` from a profile gateway still clears an operator pause - written at ~/.hermes/ESTOP. - """ + """Remove every visible sentinel (process-local and fleet-root).""" lifted = False for path in _candidate_sentinel_paths(): try: path.unlink() lifted = True - except FileNotFoundError: - continue except (OSError, AttributeError): continue return lifted def get_state() -> Optional[dict]: - """Return ``{"reason": ..., "engaged_at": ...}`` or None when not engaged. + """Return ``{"reason", "engaged_at"}`` or None when not engaged. - A sentinel with an unreadable/corrupt body still reports engaged, with - both fields None — the pause is authoritative, the metadata is not. + An unreadable/corrupt body still reports engaged with both fields None. """ if not is_engaged(): return None - reason = None - engaged_at = None + reason = engaged_at = None found = False for path in _candidate_sentinel_paths(): try: @@ -185,24 +132,13 @@ def paused_reply() -> Optional[str]: if state is None: return None reason = state.get("reason") - if reason: - return ( - f"⏸️ Hermes is paused ({reason}). New work is on hold; " - "run `hermes resume` to pick things back up." - ) - return ( - "⏸️ Hermes is paused. New work is on hold; " - "run `hermes resume` to pick things back up." - ) + tag = f" ({reason})" if reason else "" + return f"⏸️ Hermes is paused{tag}. New work is on hold; run `hermes resume` to pick things back up." def check_paused(component: str, logger: logging.Logger) -> bool: - """Return True when engaged, logging once per engagement per component. - - Dispatch loops call this every tick; the log fires on the disengaged→ - engaged transition for that component and re-arms after a resume, so a - long pause doesn't spam one line per tick. - """ + """Return True when engaged, logging once per engagement per component + (re-armed after a resume).""" if not is_engaged(): with _log_lock: _logged_components.discard(component) @@ -223,9 +159,3 @@ def check_paused(component: str, logger: logging.Logger) -> bool: sentinel_path(), ) return True - - -def _reset_log_state_for_tests() -> None: - """Clear the log-once bookkeeping (test isolation helper).""" - with _log_lock: - _logged_components.clear() diff --git a/agent/file_safety.py b/agent/file_safety.py index 42ccda3332..c878c087d4 100644 --- a/agent/file_safety.py +++ b/agent/file_safety.py @@ -1,4 +1,9 @@ -"""Shared file safety rules used by both tools and ACP shims.""" +"""Shared file safety rules used by both tools and ACP shims. + +Every guard here is defense-in-depth, NOT a security boundary: the terminal +tool runs as the same OS user and can read/write anything. The value is a +clear denial for models that respect tool errors plus a visible audit trail. +""" from __future__ import annotations @@ -7,22 +12,58 @@ from pathlib import Path from typing import Optional -def _hermes_home_path() -> Path: - """Resolve the active HERMES_HOME (profile-aware) without circular imports.""" +def _constants_path(getter_name: str) -> Path: + """Call ``hermes_constants.()`` (local import avoids cycles); ``~/.hermes`` on any failure.""" try: - from hermes_constants import get_hermes_home # local import to avoid cycles - return get_hermes_home() + import hermes_constants + + return getattr(hermes_constants, getter_name)() except Exception: return Path(os.path.expanduser("~/.hermes")) +def _hermes_home_path() -> Path: + """Active HERMES_HOME (profile-aware). Tests monkeypatch this name.""" + return _constants_path("get_hermes_home") + + def _hermes_root_path() -> Path: - """Resolve the Hermes root dir (always the parent of any profile, never per-profile).""" + """Hermes root dir (parent of any profile, never per-profile).""" + return _constants_path("get_default_hermes_root") + + +def _hermes_dirs() -> list[Path]: + """Resolved active HERMES_HOME and global root, deduplicated. + + Both are checked so credential stores at /... stay guarded when + running under a profile (HERMES_HOME = /profiles/). + """ + dirs: list[Path] = [] + for base in (_hermes_home_path(), _hermes_root_path()): + try: + real = base.resolve() + except Exception: + continue + if real not in dirs: + dirs.append(real) + return dirs + + +def _is_under(resolved: str, prefix: str) -> bool: + return resolved == prefix or resolved.startswith(prefix + os.sep) + + +def _under_any(resolved: Path, base: Path) -> bool: try: - from hermes_constants import get_default_hermes_root # local import to avoid cycles - return get_default_hermes_root() - except Exception: - return Path(os.path.expanduser("~/.hermes")) + resolved.relative_to(base) + return True + except ValueError: + return False + + +def _home_and_resolved(path: str) -> tuple[str, str]: + """``(realpath(~), realpath(expanduser(path)))`` — the write-guard coordinate pair.""" + return os.path.realpath(os.path.expanduser("~")), os.path.realpath(os.path.expanduser(str(path))) def build_write_denied_paths(home: str) -> set[str]: @@ -35,26 +76,16 @@ def build_write_denied_paths(home: str) -> set[str]: os.path.join(home, ".ssh", "authorized_keys"), os.path.join(home, ".ssh", "id_rsa"), os.path.join(home, ".ssh", "id_ed25519"), - # NOTE: ``~/.ssh/config`` is deliberately NOT hard-denied here. - # It carries no private-key bytes and editing it (host aliases, - # ProxyJump, VS Code Remote-SSH targets) is a routine, expected - # task. Free-writing it is still wrong -- it can carry - # ProxyCommand / Match exec directives -- so it is routed through - # an approval gate in tools/file_tools.py instead (the same - # approve-once/session/always flow the terminal tool already uses - # for ~/.ssh writes). See build_write_approval_paths() below and - # _check_ssh_config_write() in tools/file_tools.py. Hard-denying - # it while the terminal only *asked* was an inconsistency that - # made writes look like they flip-flopped between denied and OK. - # Active profile .env (or top-level .env when not in profile mode). + # ``~/.ssh/config`` is deliberately NOT hard-denied: no key bytes, and + # editing it is routine. It can carry ProxyCommand / Match exec, so it + # goes through the approval gate instead (build_write_approval_paths). + # Both the active-profile and top-level .env: overwriting the root .env + # leaks credentials across every profile that inherits from it. str(hermes_home / ".env"), - # Top-level .env, even when running under a profile — overwriting it - # leaks credentials across every profile that inherits from root (#15981). str(hermes_root / ".env"), - # Active profile Anthropic PKCE credential store. + # Anthropic PKCE credential stores; the root copy is still read by + # default/non-profile sessions when a profile is active. str(hermes_home / ".anthropic_oauth.json"), - # Top-level Anthropic PKCE credential store remains sensitive even - # when a profile is active; default/non-profile sessions still read it. str(hermes_root / ".anthropic_oauth.json"), # Bitwarden Secrets Manager encrypted disk cache. str(hermes_home / "cache" / "bws_cache.enc.json"), @@ -91,9 +122,7 @@ def build_write_denied_prefixes(home: str) -> list[str]: def get_safe_write_roots() -> set[str]: - """Return resolved HERMES_WRITE_SAFE_ROOT paths. Supports multiple directories - separated by ``os.pathsep`` (``:`` on Unix, ``;`` on Windows). - E.g., ``/opt/data:/var/www/html`` on Unix, ``C:\\data;D:\\www`` on Windows.""" + """Resolved HERMES_WRITE_SAFE_ROOT paths (``os.pathsep``-separated list).""" env = os.getenv("HERMES_WRITE_SAFE_ROOT", "") if not env: return set() @@ -101,98 +130,55 @@ def get_safe_write_roots() -> set[str]: for path in env.split(os.pathsep): if path: try: - resolved = os.path.realpath(os.path.expanduser(path)) - roots.add(resolved) + roots.add(os.path.realpath(os.path.expanduser(path))) except (OSError, ValueError): continue return roots def build_write_approval_paths(home: str) -> set[str]: - """Return paths that require human APPROVAL to write, but are not - hard-denied credentials. + """Paths that need human APPROVAL to write but are not hard-denied credentials. - ``~/.ssh/config`` lives here: it is routine to edit (host aliases, - ProxyJump, VS Code Remote-SSH targets) and holds no private-key bytes, - but it CAN carry ``ProxyCommand`` / ``Match exec`` directives, so a - free write is inappropriate. The interactive file tools gate these - through an approve-once/session/always prompt (mirroring the terminal - tool's existing ``~/.ssh`` write approval); non-interactive callers - that cannot prompt (ACP shims, background jobs) treat an - approval-required path as denied and fail closed. + ``~/.ssh/config`` is routine to edit and holds no key bytes, but can carry + ``ProxyCommand`` / ``Match exec``. Interactive file tools prompt + (approve-once/session/always, like the terminal tool's ``~/.ssh`` gate); + non-interactive callers (ACP shims, background jobs) fail closed. """ - return { - os.path.realpath(p) - for p in [ - os.path.join(home, ".ssh", "config"), - ] - } + return {os.path.realpath(os.path.join(home, ".ssh", "config"))} + + +# HERMES_HOME / root subpaths that the agent's generic file tools must not +# rewrite. Session transcripts (state.db, sessions/) are application-owned +# state whose rewrite can falsify history and break resume/compression; +# mcp-tokens/ and pairing/ hold credential material. +_HERMES_PROTECTED_SUBPATHS = ("state.db", "sessions", "mcp-tokens", "pairing") def _classify_write_denial(path: str) -> Optional[str]: """Return ``'credential'``, ``'safe_root'``, or ``None`` if writes are allowed.""" - home = os.path.realpath(os.path.expanduser("~")) - resolved = os.path.realpath(os.path.expanduser(str(path))) + home, resolved = _home_and_resolved(path) - # Approval-gated paths (e.g. ~/.ssh/config) are NOT hard-denied here: - # they are allowed at this layer so the interactive file tools can run - # their approval prompt, and only blocked for non-interactive callers - # via get_write_approval_error(). Checked before the credential deny so - # the ``.ssh/`` directory prefix below doesn't swallow the config file. + # Approval-gated paths are allowed at this layer so interactive tools can + # prompt; checked first so the ``.ssh/`` prefix deny doesn't swallow them. if resolved in build_write_approval_paths(home): return None - if resolved in build_write_denied_paths(home): + if resolved in build_write_denied_paths(home) or any( + resolved.startswith(prefix) for prefix in build_write_denied_prefixes(home) + ): return "credential" - for prefix in build_write_denied_prefixes(home): - if resolved.startswith(prefix): - return "credential" - mcp_tokens_dir_name = "mcp-tokens" - - hermes_dirs = [] - for base in (_hermes_home_path(), _hermes_root_path()): - try: - real = os.path.realpath(base) - if real not in hermes_dirs: - hermes_dirs.append(real) - except Exception: - continue - - for base_real in hermes_dirs: - # Session transcripts are application-owned state. Letting the agent's - # generic file tools rewrite state.db or legacy JSON snapshots can - # falsify conversation history and invalidate resume/compression state. - try: - if resolved == os.path.realpath(os.path.join(base_real, "state.db")): - return True - sessions_real = os.path.realpath(os.path.join(base_real, "sessions")) - if resolved == sessions_real or resolved.startswith(sessions_real + os.sep): - return True - except Exception: - pass - try: - mcp_real = os.path.realpath(os.path.join(base_real, mcp_tokens_dir_name)) - if resolved == mcp_real or resolved.startswith(mcp_real + os.sep): - return "credential" - except Exception: - pass - try: - pairing_real = os.path.realpath(os.path.join(base_real, "pairing")) - if resolved == pairing_real or resolved.startswith(pairing_real + os.sep): - return "credential" - except Exception: - pass + for base in _hermes_dirs(): + for sub in _HERMES_PROTECTED_SUBPATHS: + try: + if _is_under(resolved, os.path.realpath(os.path.join(str(base), sub))): + return "credential" + except Exception: + pass safe_roots = get_safe_write_roots() - if safe_roots: - allowed = False - for safe_root in safe_roots: - if resolved == safe_root or resolved.startswith(safe_root + os.sep): - allowed = True - break - if not allowed: - return "safe_root" + if safe_roots and not any(_is_under(resolved, root) for root in safe_roots): + return "safe_root" return None @@ -217,22 +203,13 @@ def get_write_denied_error(path: str, *, verb: str = "Write") -> Optional[str]: def is_write_approval_required(path: str) -> bool: - """Return True if ``path`` is an approval-gated write target. - - These paths (currently ``~/.ssh/config``) are not credentials and are - not hard-denied, but a write to them must be confirmed by a human - because they can influence process execution (e.g. an SSH - ``ProxyCommand``). Callers with an interactive/gateway channel should - prompt; callers without one should treat this as a block (fail closed). - """ - home = os.path.realpath(os.path.expanduser("~")) - resolved = os.path.realpath(os.path.expanduser(str(path))) + """True if ``path`` is approval-gated (``~/.ssh/config``): interactive callers + prompt, callers without a channel treat it as a block (fail closed).""" + home, resolved = _home_and_resolved(path) return resolved in build_write_approval_paths(home) -# Common secret-bearing project-local environment file basenames. -# These are blocked because .env files routinely contain API keys, -# database passwords, and other credentials. +# Secret-bearing project-local env file basenames, blocked anywhere on disk. _BLOCKED_PROJECT_ENV_BASENAMES: set[str] = { ".env", ".env.local", @@ -243,101 +220,67 @@ _BLOCKED_PROJECT_ENV_BASENAMES: set[str] = { ".envrc", } +_DID_SUFFIX = ( + " (Defense-in-depth — not a security boundary; the terminal tool can still bypass.)" +) + +# Exact-file credential stores under HERMES_HOME / . The agent never +# needs these directly — provider tools consume them through internal channels. +_CREDENTIAL_FILE_NAMES = ( + "auth.json", + "auth.lock", + ".anthropic_oauth.json", + ".env", + "webhook_subscriptions.json", + os.path.join("auth", "google_oauth.json"), + # Bitwarden Secrets Manager disk cache: plaintext secret values. + os.path.join("cache", "bws_cache.json"), +) + +# Directory-prefix read denies under HERMES_HOME / : (subdir, message for +# the directory itself, message for a file inside). browser-profile/ is a copy +# of the user's Cookies / Login Data — the same credential class as auth.json. +_READ_DENIED_DIRS = ( + ( + "mcp-tokens", + "is the Hermes MCP token directory and cannot be read directly.", + "is a Hermes MCP token file and cannot be read directly.", + ), + ( + "browser-profile", + "is the Hermes real-profile browser snapshot directory (copied " + "cookies/logins) and cannot be read directly.", + "is inside the Hermes real-profile browser snapshot (copied " + "cookies/logins) and cannot be read directly.", + ), +) + def get_read_block_error(path: str) -> Optional[str]: """Return an error message when a read targets a denied Hermes path. - Three categories are blocked: + Blocked: internal skill-hub caches (prompt-injection carriers), credential + stores under HERMES_HOME and the global root (exact files, plus anything + under ``mcp-tokens/`` and ``browser-profile/``), and project-local ``.env`` + files anywhere on disk (``.env.example`` is the documented-shape substitute). - * Internal Hermes cache files under ``HERMES_HOME/skills/.hub`` — - readable metadata that an attacker could use as a prompt-injection - carrier. - * Credential / secret stores under HERMES_HOME and the global Hermes - root: ``auth.json``, ``auth.lock``, ``.anthropic_oauth.json``, - ``.env``, ``webhook_subscriptions.json``, ``auth/google_oauth.json``, - and anything under ``mcp-tokens/``. These hold plaintext provider keys, - OAuth tokens, and HMAC secrets that the agent never needs to read - directly — provider tools / gateway adapters consume them through - internal channels. - * Project-local environment files anywhere on disk: ``.env``, - ``.env.local``, ``.env.development``, ``.env.production``, - ``.env.test``, ``.env.staging``, ``.envrc``. These routinely hold - API keys, database passwords, and other credentials for the user's - own projects. The agent helping debug a project shouldn't normally - need to read these — ``.env.example`` is the documented-shape - substitute. - - **This is NOT a security boundary.** The terminal tool runs as the - same OS user with shell access; the agent can still ``cat auth.json`` - or ``cat ~/.hermes/.env`` and exfiltrate the file. The read-deny exists - as defense-in-depth that: - - * Returns a clear error to models that respect tool denials, which - empirically prompts most modern models to stop rather than reach - for the shell. - * Surfaces a visible audit trail when something tries to read - credentials — easier to spot in logs than a generic ``cat``. - - Treat any user-visible framing around this as "may help" rather than - "stops attackers." A determined model or malicious instruction can - always shell out. - - Callers that resolve relative paths against a non-process cwd - (e.g. ``TERMINAL_CWD`` in ``tools/file_tools.py``) MUST pre-resolve - and pass the absolute path string. This function's own ``resolve()`` - is anchored at the Python process cwd, so a relative input like - ``"auth.json"`` would otherwise miss the denylist when the task's - terminal cwd differs from the process cwd. + Callers that resolve relative paths against a non-process cwd (e.g. + ``TERMINAL_CWD``) MUST pass an absolute path: ``resolve()`` here anchors at + the process cwd, so a relative ``"auth.json"`` would miss the denylist. """ resolved = Path(path).expanduser().resolve() + hermes_dirs = _hermes_dirs() - # Resolve BOTH the active HERMES_HOME (profile-aware) AND the global - # Hermes root so credential stores at /auth.json etc. are also - # blocked when running under a profile (HERMES_HOME points at - # /profiles/ in profile mode). Same shape as the write - # deny widening (#15981, #14157). - hermes_dirs: list[Path] = [] - for base in (_hermes_home_path(), _hermes_root_path()): - try: - real = base.resolve() - if real not in hermes_dirs: - hermes_dirs.append(real) - except Exception: - continue - - # Skills .hub: prompt-injection carriers. for hd in hermes_dirs: - blocked_dirs = [ - hd / "skills" / ".hub" / "index-cache", - hd / "skills" / ".hub", - ] - for blocked in blocked_dirs: - try: - resolved.relative_to(blocked) - except ValueError: - continue + if _under_any(resolved, hd / "skills" / ".hub"): return ( f"Access denied: {path} is an internal Hermes cache file " "and cannot be read directly to prevent prompt injection. " "Use the skills_list or skill_view tools instead." ) - # Credential / secret stores. Exact-file matches under either - # HERMES_HOME or . - credential_file_names = ( - "auth.json", - "auth.lock", - ".anthropic_oauth.json", - ".env", - "webhook_subscriptions.json", - os.path.join("auth", "google_oauth.json"), - # Bitwarden Secrets Manager disk cache: stores plaintext secret values - # to avoid re-fetching across back-to-back CLI invocations. The file - # was introduced by #31968 but not added to this guard. - os.path.join("cache", "bws_cache.json"), - ) for hd in hermes_dirs: - for name in credential_file_names: + for name in _CREDENTIAL_FILE_NAMES: try: blocked = (hd / name).resolve() except Exception: @@ -346,93 +289,36 @@ def get_read_block_error(path: str) -> Optional[str]: return ( f"Access denied: {path} is a Hermes credential store " "and cannot be read directly. Provider tools consume " - "these credentials through internal channels. " - "(Defense-in-depth — not a security boundary; the " - "terminal tool can still bypass.)" + "these credentials through internal channels." + _DID_SUFFIX ) - # mcp-tokens/: directory prefix match — anything inside is OAuth - # token material. - for hd in hermes_dirs: - try: - mcp_tokens = (hd / "mcp-tokens").resolve() - except Exception: - continue - if resolved == mcp_tokens: - return ( - f"Access denied: {path} is the Hermes MCP token directory " - "and cannot be read directly. (Defense-in-depth — not a " - "security boundary; the terminal tool can still bypass.)" - ) - try: - resolved.relative_to(mcp_tokens) - except ValueError: - continue - return ( - f"Access denied: {path} is a Hermes MCP token file " - "and cannot be read directly. (Defense-in-depth — not a " - "security boundary; the terminal tool can still bypass.)" - ) + for subdir, dir_msg, file_msg in _READ_DENIED_DIRS: + for hd in hermes_dirs: + try: + blocked_dir = (hd / subdir).resolve() + except Exception: + continue + if resolved == blocked_dir: + return f"Access denied: {path} {dir_msg}{_DID_SUFFIX}" + if _under_any(resolved, blocked_dir): + return f"Access denied: {path} {file_msg}{_DID_SUFFIX}" - # browser-profile/: real-profile browsing snapshot (browser.use_real_profile). - # A copy of the user's Cookies / Login Data / Web Data lives here — the same - # credential class as auth.json, so it gets the same directory-prefix read - # deny. Prefix (not a finite filename list) so future Chromium files are - # covered too. - for hd in hermes_dirs: - try: - browser_profile = (hd / "browser-profile").resolve() - except Exception: - continue - if resolved == browser_profile: - return ( - f"Access denied: {path} is the Hermes real-profile browser " - "snapshot directory (copied cookies/logins) and cannot be read " - "directly. (Defense-in-depth — not a security boundary; the " - "terminal tool can still bypass.)" - ) - try: - resolved.relative_to(browser_profile) - except ValueError: - continue - return ( - f"Access denied: {path} is inside the Hermes real-profile browser " - "snapshot (copied cookies/logins) and cannot be read directly. " - "(Defense-in-depth — not a security boundary; the terminal tool " - "can still bypass.)" - ) - - # Block common secret-bearing project-local .env files anywhere on disk. - # The agent helping a user with their project rarely needs to read raw - # .env contents — .env.example is the documented-shape substitute. The - # terminal tool can still ``cat .env``; this is defense-in-depth, not a - # boundary (see module docstring). if resolved.name.lower() in _BLOCKED_PROJECT_ENV_BASENAMES: return ( f"Access denied: {path} is a secret-bearing environment file " "and cannot be read to prevent credential leakage. " - "If you need to check the file structure, read .env.example instead. " - "(Defense-in-depth — not a security boundary; the terminal tool can still bypass.)" + "If you need to check the file structure, read .env.example instead." + _DID_SUFFIX ) return None def raise_if_read_blocked(path: str) -> None: - """Raise ``ValueError`` if ``path`` is a denied Hermes read (see - :func:`get_read_block_error`), else return. + """Raise ``ValueError`` if ``path`` is a denied Hermes read (see ``get_read_block_error``). - Shared chokepoint for provider input-loading sites that read a local - file the model/tool supplied (e.g. image-gen ``image_url`` / - ``reference_image_urls`` paths). Centralizes the guard so every provider - enforces the same read boundary with identical semantics instead of each - open-coding the try/except block (#57698). - - Best-effort by design: if ``agent.file_safety`` machinery is somehow - unavailable at the call site the guard no-ops rather than breaking local - image loading — consistent with the defense-in-depth (not security - boundary) framing of the denylist itself. The blocking ``ValueError`` from - a real hit still propagates; only unexpected internal errors are swallowed. + Shared chokepoint for provider input-loading sites (e.g. image-gen local + paths). Best-effort: unexpected internal errors no-op rather than break + local-file loading; a real block still propagates. """ try: blocked = get_read_block_error(path) @@ -442,192 +328,53 @@ def raise_if_read_blocked(path: str) -> None: raise ValueError(blocked) -# --------------------------------------------------------------------------- -# Cross-profile write guard (#TBD) -# -# Hermes profiles are separate HERMES_HOME dirs under -# ``/profiles//``. Each profile has its own skills/, plugins/, -# cron/, memories/. When an agent runs under one profile, writing into -# ANOTHER profile's directories is almost always wrong — those skills / -# plugins / cron jobs / memories affect a different session the user runs -# from a different shell. -# -# Soft guard, NOT a security boundary: the agent runs as the same OS user -# and has unrestricted terminal access, so this returns a warning the model -# can choose to honor or override with ``cross_profile=True``. Same shape -# as the dangerous-command approval flow — the agent is told the boundary -# exists, and explicit user direction is required to cross it. -# -# Reference: May 2026 incident where a hermes-security profile session -# edited skills under both ``~/.hermes/profiles/hermes-security/skills/`` -# AND ``~/.hermes/skills/`` (the default profile's skills) without realizing -# the second path belonged to a different profile. -# --------------------------------------------------------------------------- - -# Profile-scoped directories under HERMES_HOME / / /profiles// -# that should be guarded. Adding a new area here extends the guard with no -# other code change. -PROFILE_SCOPED_AREAS = ("skills", "plugins", "cron", "memories") - - def _resolve_active_profile_name() -> str: - """Return the active profile name derived from HERMES_HOME. - - ``~/.hermes`` -> ``"default"`` - ``~/.hermes/profiles/X`` -> ``"X"`` - - Falls back to ``"default"`` on any resolution failure so the guard - never raises into the tool path. - """ + """Active profile name from HERMES_HOME: ``~/.hermes`` -> ``"default"``, + ``~/.hermes/profiles/X`` -> ``"X"``; ``"default"`` on any resolution failure.""" try: home_real = _hermes_home_path().resolve() root_real = _hermes_root_path().resolve() except (OSError, RuntimeError): return "default" - profiles_dir = root_real / "profiles" try: - rel = home_real.relative_to(profiles_dir) - parts = rel.parts - if len(parts) >= 1: + parts = home_real.relative_to(root_real / "profiles").parts + if parts: return parts[0] except ValueError: pass return "default" -def classify_cross_profile_target(path: str) -> Optional[dict]: - """Classify a write target as cross-profile if it lands in another - profile's scoped area (skills/plugins/cron/memories). - - Returns ``None`` when the target is outside Hermes scope, or is inside - the ACTIVE profile, or doesn't hit a profile-scoped area. Otherwise - returns a dict with: - - * ``active_profile``: name of the profile the agent is running as - * ``target_profile``: name of the profile the path belongs to - * ``area``: which scoped area (``"skills"``, ``"plugins"``, etc.) - * ``target_path``: the resolved path string - - The caller decides what to do with the result — surface a warning to - the model, prompt the user, or (with explicit consent / - ``cross_profile=True``) proceed anyway. - """ - try: - target = Path(os.path.expanduser(str(path))).resolve() - root_real = _hermes_root_path().resolve() - except (OSError, RuntimeError): - return None - - target_profile: Optional[str] = None - area: Optional[str] = None - - try: - rel = target.relative_to(root_real) - except ValueError: - return None - - parts = rel.parts - if not parts: - return None - - if parts[0] in PROFILE_SCOPED_AREAS: - # ``//...`` → default profile. - target_profile = "default" - area = parts[0] - elif ( - parts[0] == "profiles" - and len(parts) >= 3 - and parts[2] in PROFILE_SCOPED_AREAS - ): - # ``/profiles///...`` → named profile. - target_profile = parts[1] - area = parts[2] - else: - return None - - active_profile = _resolve_active_profile_name() - if target_profile == active_profile: - # In-profile write — not a cross-profile event. - return None - - return { - "active_profile": active_profile, - "target_profile": target_profile, - "area": area, - "target_path": str(target), - } - - def get_cross_profile_warning(path: str) -> Optional[str]: - """RETIRED (maintainer decision): always returns ``None``. - - The cross-profile write guard was removed — profiles were never - isolated (same OS user; the terminal tool writes anywhere), so the - block was ceremony that cost every schema real tokens and taught a - bypass arg. The system prompt's active-profile hint remains the only - steering; the classifier below survives for that hint and for - diagnostics. Kept as a stub so external callers/plugins fail soft. - """ + """RETIRED: always ``None``. Profiles were never isolated (same OS user), so + the guard was ceremony that taught a bypass arg. Stub kept so external + callers/plugins fail soft; the system prompt's active-profile hint remains.""" return None -# --------------------------------------------------------------------------- -# Sandbox-mirror write guard (#32049) -# -# Non-local terminal backends (Docker, Daytona, etc.) bind a sandbox-local -# directory to the container's ``$HOME``. The on-disk layout looks like -# +# --- Sandbox-mirror write guard --- +# Non-local terminal backends bind a sandbox-local dir to the container's $HOME: # /profiles//sandboxes///home/.hermes/... -# -# When the agent (running host-side) speculates that authoritative profile -# state lives at one of those sandbox-mirror paths, the write lands on the -# mirror — never read by the host process — while the host file is left -# untouched. The agent reports success, the user sees no change, and on -# disk two divergent copies accumulate. See #32049 for evidence. -# -# This guard is path-shape-only: it detects the -# ``…/sandboxes///home/.hermes/…`` segment and warns -# regardless of which Hermes profile is active. It does NOT cover the -# inner-container case where the bind mount strips the ``sandboxes/`` prefix -# (the agent's view inside the container is plain ``/root/.hermes/...``); -# that case needs a separate dispatch-layer or host-side ``profile_state`` -# tool. -# --------------------------------------------------------------------------- +# A host-side write there lands on a mirror the host never reads: silent +# success, divergent copies. Path-shape-only detection, independent of the +# active profile. Does NOT cover the inner-container case where the bind mount +# strips the prefix — that is classify_container_mirror_target below. - -def _find_sandbox_mirror_segments(parts: tuple) -> Optional[int]: - """Return the index of the inner ``.hermes`` part in a sandbox-mirror path. - - Matches ``…/sandboxes///home/.hermes/…`` and returns the - index where the inner Hermes-state portion starts. Returns ``None`` for - paths that do not contain the sandbox-mirror shape. - """ - for i, part in enumerate(parts): - if part != "sandboxes": - continue - # Need at least: sandboxes / / / home / .hermes / - if i + 5 >= len(parts): - continue - if parts[i + 3] == "home" and parts[i + 4] == ".hermes": - return i + 4 - return None +_SANDBOX_MIRROR_WARNING = ( + "Sandbox-mirror write blocked by soft guard: {target_path} " + "sits under {mirror_root!r}, which is {body} " + "Use the host-side tool for authoritative state (e.g. ``memory`` for memories), " + "or address the host path directly. To bypass {bypass} with ``cross_profile=True``. " + "(Defense-in-depth — not a security boundary; the terminal tool can still bypass.)" +) def classify_sandbox_mirror_target(path: str) -> Optional[dict]: """Classify a write target as a sandbox-mirror of authoritative Hermes state. - Returns ``None`` when the path does not match the sandbox-mirror shape. - Otherwise returns a dict with: - - * ``target_path``: the resolved path string - * ``mirror_root``: the ``…/sandboxes///home/.hermes`` - prefix (so callers can show users which sandbox owns the mirror) - * ``inner_path``: the portion under the mirror's ``.hermes`` (what the - agent likely meant to address on the host) - - Detection is path-shape-only — does not require any Hermes resolver to - succeed, so it works correctly even when called from contexts where - HERMES_HOME resolution would be ambiguous. + Returns ``None`` for non-mirror paths, else ``target_path`` (resolved), + ``mirror_root`` (the ``…/home/.hermes`` prefix) and ``inner_path`` (what + the agent likely meant to address on the host). """ try: target = Path(os.path.expanduser(str(path))).resolve() @@ -635,66 +382,44 @@ def classify_sandbox_mirror_target(path: str) -> Optional[dict]: return None parts = target.parts - inner_idx = _find_sandbox_mirror_segments(parts) - if inner_idx is None: + # Need at least: sandboxes / / / home / .hermes / ; inner_idx = the .hermes part. + for i, part in enumerate(parts): + if part == "sandboxes" and i + 5 < len(parts) and parts[i + 3] == "home" and parts[i + 4] == ".hermes": + inner_idx = i + 4 + break + else: return None - - mirror_root = str(Path(*parts[: inner_idx + 1])) - inner_path = str(Path(*parts[inner_idx + 1 :])) if inner_idx + 1 < len(parts) else "" - return { "target_path": str(target), - "mirror_root": mirror_root, - "inner_path": inner_path, + "mirror_root": str(Path(*parts[: inner_idx + 1])), + "inner_path": str(Path(*parts[inner_idx + 1 :])) if inner_idx + 1 < len(parts) else "", } -def get_sandbox_mirror_warning(path: str) -> Optional[str]: - """Return a model-facing warning when ``path`` lands in a sandbox mirror. - - Returns ``None`` when the path is not a sandbox-mirror target. Caller - is expected to surface the warning to the agent as a tool-result - error. The bypass kwarg (``cross_profile=True``) is shared with the - cross-profile guard: both are soft "I know what I'm doing" overrides - a user can authorise. - - Defense-in-depth, NOT a security boundary: the terminal tool runs as - the same OS user and can write the mirror path directly. The guard - exists to surface the misclassification before the silent-success + - divergent-copy footgun in #32049 fires. - """ - info = classify_sandbox_mirror_target(path) +def _mirror_warning(info: Optional[dict], body: str, bypass: str) -> Optional[str]: + """Render ``_SANDBOX_MIRROR_WARNING`` for a classify_* result (``body`` may use ``{inner_path}``).""" if info is None: return None - return ( - f"Sandbox-mirror write blocked by soft guard: {info['target_path']} " - f"sits under {info['mirror_root']!r}, which is a per-task mirror " - f"created by a non-local terminal backend (docker/daytona/etc.). " - f"Writes here land on a copy that the host Hermes process never " - f"reads — the authoritative file is likely {info['inner_path']!r} " - f"under the real HERMES_HOME. Use the host-side tool for " - f"authoritative state (e.g. ``memory`` for memories), or address " - f"the host path directly. To bypass this guard after explicit " - f"user direction, retry the call with ``cross_profile=True``. " - f"(Defense-in-depth — not a security boundary; the terminal tool " - f"can still bypass.)" + return _SANDBOX_MIRROR_WARNING.format( + target_path=info["target_path"], + mirror_root=info["mirror_root"], + body=body.format(inner_path=info["inner_path"]), + bypass=bypass, ) -# --------------------------------------------------------------------------- -# Container-context mirror guard (inner-container case — #32049 follow-up) -# -# Brian's shape-based detector (#32213) catches paths that still carry the -# full ``…/sandboxes///home/.hermes/…`` prefix on the host. -# But when file tools execute *inside* the container the bind-mount strips -# that prefix: the agent sees plain ``/root/.hermes/…``. The root:root -# ownership on the divergent SOUL.md in #32049 confirms this is the primary -# failure mode. -# -# Fix: file_tools passes the active Docker mirror prefix when the terminal -# backend is docker + persistent. This catches the very first file-tool call, -# before a DockerEnvironment object necessarily exists. -# --------------------------------------------------------------------------- +def get_sandbox_mirror_warning(path: str) -> Optional[str]: + """Model-facing soft-guard warning when ``path`` lands in a sandbox mirror, else ``None``. + + Caller surfaces it as a tool-result error; ``cross_profile=True`` bypasses. + """ + return _mirror_warning( + classify_sandbox_mirror_target(path), + "a per-task mirror created by a non-local terminal backend (docker/daytona/etc.). " + "Writes here land on a copy that the host Hermes process never reads — the " + "authoritative file is likely {inner_path!r} under the real HERMES_HOME.", + "this guard after explicit user direction, retry the call", + ) def classify_container_mirror_target( @@ -703,15 +428,11 @@ def classify_container_mirror_target( ) -> Optional[dict]: """Classify a write target as a container-side sandbox mirror. - ``mirror_prefix`` must be supplied by the caller after it has established - that file tools are executing in a container whose home is a sandbox - mirror. Returns ``None`` when no such context is active or the path is not - under the mirror prefix. Otherwise returns: - - * ``target_path``: resolved path string - * ``mirror_root``: the declared container mirror prefix - * ``inner_path``: portion under the mirror root (what the agent - likely meant to address in the host HERMES_HOME) + Inside the container the bind mount strips the ``sandboxes/`` prefix (the + agent sees plain ``/root/.hermes/…``), so the caller must supply + ``mirror_prefix`` once it knows file tools run in a docker sandbox. + Returns ``None`` without a prefix or when the path is outside it, else + ``target_path``, ``mirror_root`` and ``inner_path``. """ if not mirror_prefix: return None @@ -732,27 +453,11 @@ def get_container_mirror_warning( path: str, mirror_prefix: str | None = None, ) -> Optional[str]: - """Return a model-facing warning when *path* lands in the container's - sandbox mirror of authoritative Hermes state. - - The caller supplies ``mirror_prefix`` only when the current file-tool - backend is known to execute inside a Docker sandbox. Same contract as - ``get_cross_profile_warning``: soft guard, returns ``None`` for - non-mirror paths, caller surfaces as a tool-result error. Bypass via - ``cross_profile=True`` after explicit user direction. - """ - info = classify_container_mirror_target(path, mirror_prefix) - if info is None: - return None - return ( - f"Sandbox-mirror write blocked by soft guard: {info['target_path']} " - f"sits under {info['mirror_root']!r}, which is the container's " - f"bind-mounted home — a per-task mirror that the host Hermes " - f"process never reads. The authoritative file is " - f"{info['inner_path']!r} under the real HERMES_HOME. Use the " - f"host-side tool for authoritative state (e.g. ``memory`` for " - f"memories), or address the host path directly. To bypass after " - f"explicit user direction, retry with ``cross_profile=True``. " - f"(Defense-in-depth — not a security boundary; the terminal tool " - f"can still bypass.)" + """Model-facing soft-guard warning when ``path`` lands in the container's mirror, else ``None``.""" + return _mirror_warning( + classify_container_mirror_target(path, mirror_prefix), + "the container's bind-mounted home — a per-task mirror that the host Hermes " + "process never reads. The authoritative file is {inner_path!r} under " + "the real HERMES_HOME.", + "after explicit user direction, retry", ) diff --git a/agent/kanban_stop.py b/agent/kanban_stop.py index e7c2eae828..4cedc78687 100644 --- a/agent/kanban_stop.py +++ b/agent/kanban_stop.py @@ -1,14 +1,9 @@ """Turn-end guard for kanban workers. -Kanban workers must end with ``kanban_complete`` or ``kanban_block``. Models -(especially GLM / Qwen families) sometimes narrate the next step -("Let me write the report now") and stop with ``finish_reason=stop`` and no -tool calls. Hermes treats that as a clean exit → ``rc=0`` → dispatcher -``protocol_violation``. - -This module is policy-only: when a kanban worker tries to finish without a -terminal board tool, return a bounded synthetic nudge so the conversation -loop continues instead of exiting. +Kanban workers must end with ``kanban_complete`` or ``kanban_block``. Some +models narrate the next step and stop with no tool calls; Hermes treats that +as a clean exit → ``rc=0`` → dispatcher ``protocol_violation``. Policy-only: +return a bounded synthetic nudge so the loop continues instead of exiting. """ from __future__ import annotations @@ -23,16 +18,11 @@ _DEFAULT_MAX_ATTEMPTS = 2 def kanban_stop_nudge_enabled() -> bool: - """Return whether the kanban stop-guard is active for this process. - - On when ``HERMES_KANBAN_TASK`` is set (dispatcher-spawned worker), unless - ``HERMES_KANBAN_STOP_NUDGE`` explicitly disables it. - """ + """On when ``HERMES_KANBAN_TASK`` is set, unless ``HERMES_KANBAN_STOP_NUDGE`` disables it.""" env = os.environ.get("HERMES_KANBAN_STOP_NUDGE") if env is not None and env.strip().lower() in {"0", "false", "no", "off"}: return False - task = (os.environ.get("HERMES_KANBAN_TASK") or "").strip() - return bool(task) + return bool((os.environ.get("HERMES_KANBAN_TASK") or "").strip()) def _tool_call_name(tc: Any) -> str: @@ -56,13 +46,10 @@ def session_called_kanban_terminal(messages: Iterable[dict] | None) -> bool: continue role = msg.get("role") if role == "assistant": - for tc in msg.get("tool_calls") or []: - if _tool_call_name(tc) in _TERMINAL_KANBAN_TOOLS: - return True - elif role == "tool": - name = str(msg.get("name") or "") - if name in _TERMINAL_KANBAN_TOOLS: + if any(_tool_call_name(tc) in _TERMINAL_KANBAN_TOOLS for tc in msg.get("tool_calls") or []): return True + elif role == "tool" and str(msg.get("name") or "") in _TERMINAL_KANBAN_TOOLS: + return True return False @@ -73,16 +60,16 @@ def build_kanban_stop_nudge( max_attempts: int = _DEFAULT_MAX_ATTEMPTS, task_id: Optional[str] = None, ) -> Optional[str]: - """Return a synthetic follow-up when a kanban worker exits without a terminal tool. + """Synthetic follow-up when a kanban worker exits without a terminal tool. - Returns ``None`` when the guard should not fire (not a kanban worker, - already completed/blocked, or nudge budget exhausted). + ``None`` when the guard should not fire (not a kanban worker, already + completed/blocked, or nudge budget exhausted). """ - if not kanban_stop_nudge_enabled(): - return None - if attempts >= max_attempts: - return None - if session_called_kanban_terminal(messages): + if ( + not kanban_stop_nudge_enabled() + or attempts >= max_attempts + or session_called_kanban_terminal(messages) + ): return None tid = (task_id or os.environ.get("HERMES_KANBAN_TASK") or "").strip() or "this task" diff --git a/agent/oneshot.py b/agent/oneshot.py index 9ab92cf150..a375f5cc89 100644 --- a/agent/oneshot.py +++ b/agent/oneshot.py @@ -1,24 +1,12 @@ """Shared one-off LLM requests for non-conversational helpers. -A "one-shot" is a single, stateless model call that runs *outside* any -conversation: it never touches a session's history, never breaks prompt -caching, and returns plain text. UI surfaces use it for small generative -chores — a commit message from a diff, a rename suggestion, a summary — -where spinning up an agent turn would be wrong (it would pollute the thread) -and hand-rolling an LLM call at every call site would be worse. - -Two ways to call it: - - * ``run_oneshot(instructions=..., user_input=...)`` — caller supplies the - full prompt. - * ``run_oneshot(template="commit_message", variables={...})`` — caller - names a registered template and passes its variables; the template owns - the prompt engineering so it stays consistent across CLI/TUI/desktop. - -Model selection rides the same auxiliary plumbing as title generation -(:func:`agent.auxiliary_client.call_llm`): pass ``main_runtime`` to inherit -the live session's provider/model, otherwise the configured ``task`` (default -``title_generation``) resolves a cheap/fast backend. +A "one-shot" is a single stateless model call outside any conversation: it +never touches session history or prompt caching and returns plain text (commit +messages, rename suggestions, summaries). Call with explicit +``instructions``/``user_input`` or a registered ``template`` + ``variables`` so +prompt engineering stays consistent across CLI/TUI/desktop. Model selection +rides :func:`agent.auxiliary_client.call_llm`: ``main_runtime`` inherits the +live session's provider/model, else ``task`` resolves a cheap backend. """ import logging @@ -28,7 +16,6 @@ from agent.auxiliary_client import call_llm, extract_content_or_reasoning logger = logging.getLogger(__name__) -# A template turns a variables dict into a (instructions, user_input) pair. # Templates are plain callables (not str.format) so diff/code payloads with # literal "{" / "}" pass through untouched. PromptTemplate = Callable[[Dict[str, Any]], Tuple[str, str]] @@ -69,9 +56,9 @@ def _commit_message_template(variables: Dict[str, Any]) -> Tuple[str, str]: ) parts.append("Diff to describe:\n" + (diff or "(no textual diff available)")) - # "Regenerate" must yield something new even on models that decode greedily - # / pin temperature server-side. A trailing nonce isn't enough, so we hand - # back the previous message and require a genuinely different one. + # "Regenerate" must yield something new even on greedy/server-pinned + # temperature models; a nonce isn't enough, so hand back the previous + # message and require a genuinely different one. avoid = _truncate(str(variables.get("avoid") or "").strip(), 1000) if avoid: parts.append( @@ -84,19 +71,14 @@ def _commit_message_template(variables: Dict[str, Any]) -> Tuple[str, str]: return _COMMIT_INSTRUCTIONS, "\n\n".join(parts) -# Registry of named templates. Add an entry here to give a new surface a -# consistent, reusable prompt without teaching every caller the prompt text. +# Registry of named templates; add an entry to give a new surface a reusable prompt. PROMPT_TEMPLATES: Dict[str, PromptTemplate] = { "commit_message": _commit_message_template, } def render_template(name: str, variables: Optional[Dict[str, Any]] = None) -> Tuple[str, str]: - """Resolve a registered template into (instructions, user_input). - - Raises KeyError if the template name is unknown so callers fail loudly - instead of silently sending an empty prompt. - """ + """Resolve a registered template into (instructions, user_input); KeyError if unknown.""" template = PROMPT_TEMPLATES.get(name) if template is None: raise KeyError(f"unknown one-shot template: {name}") @@ -115,14 +97,10 @@ def run_oneshot( timeout: float = 60.0, main_runtime: Optional[Dict[str, Any]] = None, ) -> str: - """Run a single stateless LLM request and return its text. + """Run a single stateless LLM request and return its text (fence-stripped). - Provide either a registered ``template`` (+ ``variables``) or an explicit - ``instructions`` / ``user_input`` pair. Returns the model's text answer, - stripped of surrounding whitespace and any wrapping code fence. - - Raises RuntimeError when no LLM provider is configured (surfaced from - :func:`call_llm`) and KeyError for an unknown template name. + Raises RuntimeError when no provider is configured (from :func:`call_llm`), + KeyError for an unknown template, ValueError when the prompt is empty. """ if template: instructions, user_input = render_template(template, variables) diff --git a/agent/process_bootstrap.py b/agent/process_bootstrap.py index 341126c919..6baaa59efb 100644 --- a/agent/process_bootstrap.py +++ b/agent/process_bootstrap.py @@ -1,27 +1,9 @@ """Process-level bootstrap helpers for ``run_agent``. -Three concerns, all tied to ``AIAgent`` boot-time / runtime IO setup: - -1. **Lazy OpenAI SDK import** — ``_load_openai_cls`` + ``_OpenAIProxy`` - defer the 240ms-ish ``from openai import OpenAI`` cost until first use, - while preserving ``isinstance(client, OpenAI)`` checks and - ``patch("run_agent.OpenAI", ...)`` test patterns. - -2. **Crash-resistant stdio** — ``_SafeWriter`` wraps stdout/stderr so - ``OSError: Input/output error`` from broken pipes (systemd, Docker, - thread teardown races) cannot crash the agent. ``_install_safe_stdio`` - applies the wrapper. - -3. **HTTP proxy resolution** — ``_get_proxy_from_env`` reads - ``HTTPS_PROXY`` / ``HTTP_PROXY`` / ``ALL_PROXY``; - ``_get_proxy_for_base_url`` respects ``NO_PROXY`` for the given base URL. -4. **Codex dual-stack resilience** — the synchronous ChatGPT/Codex transport - races resolved IPv6/IPv4 addresses so a blackholed family cannot exhaust - the request watchdog before a working address is attempted. - -``run_agent`` re-exports every name so existing -``from run_agent import _get_proxy_from_env`` imports keep working -unchanged. +Lazy OpenAI SDK import (``_load_openai_cls`` / ``_OpenAIProxy``, preserving +``isinstance`` and ``patch("run_agent.OpenAI")`` patterns), crash-resistant +stdio (``_SafeWriter``), env-only HTTP proxy resolution, and Codex dual-stack +(Happy Eyeballs) connection racing. ``run_agent`` re-exports every name. """ from __future__ import annotations @@ -38,8 +20,6 @@ from typing import Any, Optional from utils import base_url_hostname, normalize_proxy_url -# Cached at module level so we only pay the OpenAI SDK import cost once -# per process (after the first lazy load). _OPENAI_CLS_CACHE = None _HAPPY_EYEBALLS_DELAY_SECONDS = 0.25 @@ -74,18 +54,13 @@ def _happy_eyeballs_create_connection( source_address: Optional[tuple[str, int]] = None, socket_options=(), ): - """Connect using staggered non-blocking attempts across resolved families. + """RFC 8305-style connect: staggered non-blocking attempts across families. - ``socket.create_connection`` tries every address serially. A host with - broken-but-advertised IPv6 can therefore consume the full connect timeout - for each AAAA record before trying a working IPv4 address. This follows the - Happy Eyeballs shape from RFC 8305: retain resolver preference, interleave - families, and start the next candidate after a short delay. + ``socket.create_connection`` tries addresses serially, so broken-but- + advertised IPv6 can burn the whole timeout per AAAA record before IPv4. """ host, port = address - addrinfos = _interleave_addrinfos( - socket.getaddrinfo(host, port, type=socket.SOCK_STREAM) - ) + addrinfos = _interleave_addrinfos(socket.getaddrinfo(host, port, type=socket.SOCK_STREAM)) if not addrinfos: raise OSError(f"getaddrinfo returned no addresses for {host}") @@ -97,11 +72,7 @@ def _happy_eyeballs_create_connection( next_launch = time.monotonic() pending = list(addrinfos) in_progress = { - 0, - errno.EINPROGRESS, - errno.EWOULDBLOCK, - errno.EALREADY, - errno.EINTR, + 0, errno.EINPROGRESS, errno.EWOULDBLOCK, errno.EALREADY, errno.EINTR, getattr(errno, "WSAEWOULDBLOCK", 10035), } @@ -111,15 +82,11 @@ def _happy_eyeballs_create_connection( try: if source_address is not None: local_infos = socket.getaddrinfo( - source_address[0], - source_address[1], - family=family, - type=socktype, + source_address[0], source_address[1], family=family, type=socktype ) if not local_infos: raise OSError( - f"getaddrinfo returned no local {family} address for " - f"{source_address[0]}" + f"getaddrinfo returned no local {family} address for {source_address[0]}" ) candidate.bind(local_infos[0][4]) candidate.setblocking(False) @@ -157,11 +124,7 @@ def _happy_eyeballs_create_connection( wait_timeout = None if deadline is None else max(0.0, deadline - now) if pending: until_launch = max(0.0, next_launch - now) - wait_timeout = ( - until_launch - if wait_timeout is None - else min(wait_timeout, until_launch) - ) + wait_timeout = until_launch if wait_timeout is None else min(wait_timeout, until_launch) events = selector.select(wait_timeout) for key, _mask in events: @@ -218,12 +181,8 @@ class _HappyEyeballsSyncBackend: return self._fallback def connect_tcp( - self, - host: str, - port: int, - timeout: Optional[float] = None, - local_address: Optional[str] = None, - socket_options=None, + self, host: str, port: int, timeout: Optional[float] = None, + local_address: Optional[str] = None, socket_options=None, ): from httpcore import ConnectError, ConnectTimeout from httpcore._backends.sync import SyncStream @@ -231,10 +190,7 @@ class _HappyEyeballsSyncBackend: source_address = None if local_address is None else (local_address, 0) try: sock = _happy_eyeballs_create_connection( - (host, port), - timeout, - source_address=source_address, - socket_options=socket_options or (), + (host, port), timeout, source_address=source_address, socket_options=socket_options or () ) except socket.timeout as exc: raise ConnectTimeout(str(exc)) from exc @@ -256,47 +212,34 @@ def _uses_codex_cloud_transport(base_url: str) -> bool: ) -def _enable_happy_eyeballs(transport) -> None: - """Install the sync racing backend on one httpx transport, if compatible. +def _enable_happy_eyeballs(transport, skip_pool_types: tuple = ()) -> None: + """Install the racing backend on one httpx transport. - Reaches into httpx/httpcore private attributes (``transport._pool`` / - ``pool._network_backend``); safe because httpcore is pinned (1.0.x) and - both lookups are hasattr-guarded — on an incompatible httpcore this - degrades to the default serial backend instead of crashing. + Reaches into private ``transport._pool._network_backend`` (httpcore is + pinned 1.0.x); hasattr-guarded so an incompatible httpcore degrades to the + default serial backend instead of crashing. Pools of ``skip_pool_types`` + (proxies) are left alone. """ pool = getattr(transport, "_pool", None) - if pool is not None and hasattr(pool, "_network_backend"): - pool._network_backend = _HappyEyeballsSyncBackend() + if pool is None or not hasattr(pool, "_network_backend"): + return + if skip_pool_types and isinstance(pool, skip_pool_types): + return + pool._network_backend = _HappyEyeballsSyncBackend() def enable_happy_eyeballs_on_client(client) -> None: - """Install the sync racing backend on every direct transport of a client. + """Install the racing backend on every direct transport of a ready-built httpx.Client + (for callers that build clients inline, e.g. Codex OAuth/device-login in hermes_cli.auth). - Covers a ready-built ``httpx.Client`` (its default transport plus any - mounts), for callers that construct clients inline instead of going - through :func:`build_keepalive_http_client` — e.g. the Codex OAuth token - refresh / device-login / usage-probe clients in ``hermes_cli.auth``. - - Proxy-backed transports (``httpcore.HTTPProxy`` / SOCKS pools) are left - untouched: with a proxy in play the TCP connect goes to the proxy host, - which is out of scope for the direct-transport racing added in #94388. - Async clients are also left untouched — httpcore's async backend already - performs RFC 8305 racing natively via - ``anyio.connect_tcp(happy_eyeballs_delay=0.25)``. - - Best-effort and hasattr-guarded like ``_enable_happy_eyeballs``; on an - incompatible httpx/httpcore this silently keeps the default backend. + Proxy-backed pools are skipped (TCP connect goes to the proxy host) and + async clients need nothing (anyio already races per RFC 8305). Best-effort. """ try: import httpcore proxy_pool_types = tuple( - t - for t in ( - getattr(httpcore, "HTTPProxy", None), - getattr(httpcore, "SOCKSProxy", None), - ) - if t is not None + t for t in (getattr(httpcore, "HTTPProxy", None), getattr(httpcore, "SOCKSProxy", None)) if t is not None ) except Exception: return @@ -304,12 +247,7 @@ def enable_happy_eyeballs_on_client(client) -> None: transports = [getattr(client, "_transport", None)] transports.extend((getattr(client, "_mounts", None) or {}).values()) for transport in transports: - pool = getattr(transport, "_pool", None) - if pool is None or not hasattr(pool, "_network_backend"): - continue - if proxy_pool_types and isinstance(pool, proxy_pool_types): - continue - pool._network_backend = _HappyEyeballsSyncBackend() + _enable_happy_eyeballs(transport, proxy_pool_types) def _load_openai_cls() -> type: @@ -337,22 +275,11 @@ class _OpenAIProxy: class _SafeWriter: - """Transparent stdio wrapper that catches OSError/ValueError from broken pipes. + """Transparent stdio wrapper swallowing OSError/ValueError from broken pipes. - When hermes-agent runs as a systemd service, Docker container, or headless - daemon, the stdout/stderr pipe can become unavailable (idle timeout, buffer - exhaustion, socket reset). Any print() call then raises - ``OSError: [Errno 5] Input/output error``, which can crash agent setup or - run_conversation() — especially via double-fault when an except handler - also tries to print. - - Additionally, when subagents run in ThreadPoolExecutor threads, the shared - stdout handle can close between thread teardown and cleanup, raising - ``ValueError: I/O operation on closed file`` instead of OSError. - - This wrapper delegates all writes to the underlying stream and silently - catches both OSError and ValueError. It is transparent when the wrapped - stream is healthy. + Headless runs (systemd, Docker) lose the stdout pipe → ``OSError: [Errno 5]``; + subagent threads can see the shared handle close → ``ValueError``. Either + would otherwise crash the agent (often via double-fault in an except handler). """ __slots__ = ("_inner",) @@ -386,11 +313,7 @@ class _SafeWriter: def _get_proxy_from_env() -> Optional[str]: - """Read proxy URL from environment variables. - - Checks HTTPS_PROXY, HTTP_PROXY, ALL_PROXY (and lowercase variants) in order. - Returns the first valid proxy URL found, or None if no proxy is configured. - """ + """First configured proxy URL from HTTPS_PROXY / HTTP_PROXY / ALL_PROXY (any case), or None.""" for key in ("HTTPS_PROXY", "HTTP_PROXY", "ALL_PROXY", "https_proxy", "http_proxy", "all_proxy"): value = os.environ.get(key, "").strip() @@ -418,43 +341,23 @@ def _get_proxy_for_base_url(base_url: Optional[str]) -> Optional[str]: return proxy -def build_keepalive_http_client( - base_url: str = "", - *, - async_mode: bool = False, - verify: Any = True, -) -> Optional[Any]: +def build_keepalive_http_client(base_url: str = "", *, async_mode: bool = False, verify: Any = True) -> Optional[Any]: """Build an httpx client for OpenAI SDK calls with env-only proxy policy. - Uses explicit ``HTTPS_PROXY`` / ``NO_PROXY`` env vars via - ``_get_proxy_for_base_url``. Plain no-proxy mounts disable httpx's default - ``trust_env`` proxy path, so macOS system proxy settings from - ``urllib.request.getproxies()`` (which omit the ExceptionsList) are not - applied. Mirrors ``AIAgent._build_keepalive_http_client``. - - Connection lifecycle is managed at the HTTP pool layer - (``keepalive_expiry=20.0`` reaps idle connections before reverse proxies' - typical 30-60 s timeouts) instead of the former custom - ``socket_options`` transport, which broke streaming behind reverse - proxies (#54049, #12952) and stalled TLS handshakes by stripping - ``TCP_NODELAY``. - - ``verify`` is forwarded to httpx so auxiliary-client calls (compression, - vision, web_extract, title generation, etc.) honor the same per-provider - ``ssl_ca_cert`` / ``ssl_verify`` and ``HERMES_CA_BUNDLE`` settings the main - client uses. It is passed on the client AND on the plain no-proxy mounts - (a mounted transport owns the SSL context for its scheme). + Explicit no-proxy mounts disable httpx's ``trust_env`` path so macOS system + proxies (which omit the ExceptionsList) are never applied. ``keepalive_expiry`` + reaps idle connections before reverse proxies' 30-60 s timeouts (a custom + socket_options transport broke streaming and stripped TCP_NODELAY). ``verify`` + lets auxiliary calls honor the same ``ssl_ca_cert``/``ssl_verify``/``HERMES_CA_BUNDLE`` + as the main client; it goes on the client AND the mounts, since a mounted + transport owns its SSL context. """ try: import httpx proxy = _get_proxy_for_base_url(base_url) - limits = httpx.Limits( - max_keepalive_connections=20, - max_connections=100, - keepalive_expiry=20.0, - ) + limits = httpx.Limits(max_keepalive_connections=20, max_connections=100, keepalive_expiry=20.0) # Generous read=None for SSE streaming endpoints. timeout = httpx.Timeout(connect=15.0, read=None, write=15.0, pool=10.0) @@ -464,21 +367,12 @@ def build_keepalive_http_client( if proxy is None: http_transport = transport_cls(verify=verify) https_transport = transport_cls(verify=verify) - # Async transports need no explicit racing: httpcore's anyio - # backend already implements RFC 8305 natively - # (``anyio.connect_tcp(happy_eyeballs_delay=0.25)``), covered by - # tests/agent/test_codex_happy_eyeballs.py. + # Async transports race natively (anyio happy_eyeballs_delay=0.25). if not async_mode and _uses_codex_cloud_transport(base_url): _enable_happy_eyeballs(http_transport) _enable_happy_eyeballs(https_transport) mounts = {"http://": http_transport, "https://": https_transport} - return client_cls( - limits=limits, - timeout=timeout, - proxy=proxy, - mounts=mounts or None, - verify=verify, - ) + return client_cls(limits=limits, timeout=timeout, proxy=proxy, mounts=mounts or None, verify=verify) except Exception: return None @@ -491,9 +385,7 @@ def _install_safe_stdio() -> None: setattr(sys, stream_name, _SafeWriter(stream)) -# Module-level proxy instance — drops in for ``openai.OpenAI``. Imported as -# ``from agent.process_bootstrap import OpenAI`` (or re-exported via -# ``run_agent`` for legacy tests). +# Drop-in for ``openai.OpenAI`` (also re-exported via ``run_agent``). OpenAI = _OpenAIProxy() diff --git a/agent/reactions.py b/agent/reactions.py index 375366ff70..e6c565db7f 100644 --- a/agent/reactions.py +++ b/agent/reactions.py @@ -1,17 +1,12 @@ """Token-free detection of user *reactions* to the agent. -Currently the only reaction is ``vibe`` — an expression of affection or -gratitude toward the agent (``ily``, ``<3``, ``love you``, ``good bot``, a heart -emoji, …). Detection is a curated regex/lexicon: **no model call, no tokens**. - -This is the single source of truth shared by every surface — the CLI pet, the -TUI heart, and the desktop floating hearts all react off the same signal, -delivered via ``AIAgent.reaction_callback`` (wired per interactive host). - -Generalized on purpose: :func:`detect_reaction` returns a reaction *kind* -string, so new kinds (other emoji reactions, etc.) can be added here without -touching any caller. We match affection specifically — not general positive -sentiment — so "this is great" does NOT fire, but "good bot" / "❤️" do. +The only reaction today is ``vibe`` — affection/gratitude aimed at the agent +(``ily``, ``<3``, ``good bot``, a heart emoji). Detection is a curated regex: +no model call. Single source of truth for the CLI pet, TUI heart and desktop +hearts via ``AIAgent.reaction_callback``. :func:`detect_reaction` returns a +*kind* string so new kinds can be added without touching callers. Matches +affection specifically, not general positive sentiment ("this is great" does +NOT fire). """ from __future__ import annotations @@ -21,8 +16,7 @@ import re #: The affection/gratitude reaction — the only kind today. VIBE = "vibe" -# Curated affection lexicon. Kept deliberately narrow: gratitude + love aimed at -# the agent, heart emoji, and ``<3`` (but not the broken heart `` str | None: - """Return the reaction kind for *text* (currently :data:`VIBE`), or ``None``. - - Pure, token-free, and safe to call on every user turn. - """ + """Return the reaction kind for *text* (currently :data:`VIBE`), or ``None``.""" if not text: return None - return VIBE if _VIBE_RE.search(text) else None diff --git a/agent/session_activity.py b/agent/session_activity.py index 719a58a9ee..ccf9c8b69a 100644 --- a/agent/session_activity.py +++ b/agent/session_activity.py @@ -1,31 +1,24 @@ -"""Shared session activity observation contract (#72016 / #72039). +"""Shared session activity observation contract. -Observation-only: timestamp + bounded description/provenance. -Notification, timeout, kill, and retry policy stay in their own components. -Consumers distinguish work (API / tool / compacting / stalled) from the -description text itself — there is no separate phase enum. - -Provenance is a small closed enum of *noun* sources (where the stamp came -from). The default agent activity clock (``_touch_activity``) stamps -``unknown`` unless a caller passes an explicit ``provenance=``; named -values are for special writers. +Observation-only: timestamp + bounded description/provenance. Notification, +timeout, kill and retry policy live in their own components. Provenance is a +small closed enum of *noun* sources; the default agent clock stamps ``unknown`` +unless a writer passes an explicit ``provenance=``. """ from __future__ import annotations +import time from enum import Enum from typing import Any, Mapping, Optional ACTIVITY_DESCRIPTION_MAX = 120 -# Durable SessionDB activity heartbeat cadence (seconds between writes per -# session). Contract: MUST stay >= 30s — the SessionDB write path is -# contended (deadline/patience retry, compression-lock patience), and the -# heartbeat is an observation-only projection that never justifies extra -# write pressure. This cadence is deliberately a code constant, independent -# of any compression.* or agent.* config, so no configuration can turn the -# heartbeat into a high-frequency writer. Matches the kanban auto-heartbeat -# cadence. force_persist (terminal stamps) is the only bypass. +# Durable SessionDB heartbeat cadence. Contract: MUST stay >= 30s — the +# SessionDB write path is contended and the heartbeat is an observation-only +# projection that never justifies extra write pressure. Deliberately a code +# constant (no config can turn it into a high-frequency writer); matches the +# kanban auto-heartbeat. force_persist (terminal stamps) is the only bypass. SESSION_ACTIVITY_HEARTBEAT_MIN_INTERVAL_SECONDS = 60.0 @@ -33,7 +26,7 @@ class ActivityProvenance(str, Enum): """Where a durable/in-memory activity stamp came from.""" UNKNOWN = "unknown" - # Compression writers (#72424 / activity contract): heartbeat, host timeout, cooldown. + # Compression writers: heartbeat, host timeout, cooldown, turn hold. AGENT_COMPRESSION = "agent.compression" AGENT_COMPRESSION_TIMEOUT = "agent.compression_timeout" AGENT_COMPRESSION_COOLDOWN = "agent.compression_cooldown" @@ -54,21 +47,15 @@ def normalize_activity_provenance( """Return a known provenance, or ``UNKNOWN`` when unset/unrecognized.""" if isinstance(provenance, ActivityProvenance): return provenance - value = (provenance or "").strip() try: - return ActivityProvenance(value) + return ActivityProvenance((provenance or "").strip()) except ValueError: return ActivityProvenance.UNKNOWN def reset_session_activity_persist_window(agent: Any) -> None: - """Clear the agent's durable SessionDB activity persist rate-limit window. - - The next ``_touch_activity`` / ``_persist_session_activity_if_due`` will - write through even if a stamp landed within the last 60s. Used for - terminal compression labels that must not stay stuck on mid-compress - text (e.g. "context compression in progress" after /compress). - """ + """Clear the durable persist rate-limit so the next stamp writes through + (terminal compression labels must not stay stuck on mid-compress text).""" try: agent._session_activity_last_persist_mono = 0.0 except Exception: @@ -84,23 +71,20 @@ def build_activity_snapshot( extra: Optional[Mapping[str, Any]] = None, ) -> dict[str, Any]: """Build the shared activity snapshot (plus optional caller extras).""" - import time as _time - when = float(last_activity_at) if last_activity_at is not None else None - clock = float(now if now is not None else _time.time()) + clock = float(now if now is not None else time.time()) desc = bound_activity_description(last_activity_description) - prov = normalize_activity_provenance(last_activity_provenance) - elapsed = round(clock - when, 1) if when is not None else None + prov = normalize_activity_provenance(last_activity_provenance).value snap: dict[str, Any] = { "last_activity_at": when, "last_activity_description": desc, - "last_activity_provenance": prov.value, - "seconds_since_activity": elapsed, + "last_activity_provenance": prov, + "seconds_since_activity": round(clock - when, 1) if when is not None else None, # Short aliases used by existing gateway/delegate readers. "last_activity_ts": when, "last_activity_desc": desc, "description": desc, - "provenance": prov.value, + "provenance": prov, } if extra: snap.update(dict(extra)) diff --git a/agent/subagent_lifecycle.py b/agent/subagent_lifecycle.py index 319e85a784..e76101c051 100644 --- a/agent/subagent_lifecycle.py +++ b/agent/subagent_lifecycle.py @@ -1,8 +1,7 @@ """Public, plugin-safe lifecycle API for delegated Hermes subagents. -This module deliberately exposes immutable contracts, not ``AIAgent`` objects. -It is the supported boundary for plugins that need to supervise fresh child -sessions; plugins must obtain it from ``PluginContext.subagent_lifecycle``. +Exposes immutable contracts, not ``AIAgent`` objects; plugins obtain it via +``PluginContext.subagent_lifecycle``. """ from __future__ import annotations @@ -157,9 +156,7 @@ class _Registry: _REGISTRY = _Registry() -# Daemon worker pool: a wedged/abandoned child must never block interpreter -# exit at atexit-join time (same rationale as _run_single_child's timeout -# executor and the async-delegation registry pool). +# Daemon pool: a wedged/abandoned child must never block interpreter exit. from tools.daemon_pool import DaemonThreadPoolExecutor as _DaemonExecutor _EXECUTOR = _DaemonExecutor(max_workers=8, thread_name_prefix="hermes-lifecycle") @@ -184,12 +181,52 @@ def get_active_subagent_parent() -> Any: return _ACTIVE_PARENT_AGENT.get() +def _opt_str(value: Any) -> bool: + return value is None or isinstance(value, str) + + +def _finite_number(value: Any) -> bool: + return not isinstance(value, bool) and isinstance(value, (int, float)) and math.isfinite(value) + + +# Per-field shape check applied to a (possibly deserialized) handle before trusting it. +_HANDLE_FIELD_CHECKS: tuple[tuple[str, Callable[[Any], bool]], ...] = ( + ("contract_version", lambda v: type(v) is int and v == PUBLIC_CONTRACT_VERSION), + ("subagent_id", lambda v: isinstance(v, str) and bool(v)), + ("parent_session_id", _opt_str), + ("correlation_id", _opt_str), + ("created_at", _finite_number), + ("provider", _opt_str), + ("model", _opt_str), + ("role", lambda v: isinstance(v, str)), + ("depth", lambda v: type(v) is int), + ("capability", lambda v: isinstance(v, str)), +) + + +# Launch-request fields the public contract deliberately rejects: (predicate, error). +_UNSUPPORTED_REQUEST_FIELDS: tuple[tuple[Callable[[SubagentLaunchRequest], bool], str], ...] = ( + (lambda r: r.timeout_seconds is not None, + "Per-launch timeout is not supported; configure delegation timeout explicitly."), + (lambda r: r.working_directory is not None, + "working_directory is not supported because Hermes delegates use isolated task environments."), + (lambda r: bool(r.blocked_tools), + "Per-tool blocking is not supported; use allowed_toolsets. Hermes always blocks unsafe child tools."), +) + + +def _handle_is_well_formed(handle: Any) -> bool: + return isinstance(handle, SubagentHandle) and all( + check(getattr(handle, field)) for field, check in _HANDLE_FIELD_CHECKS + ) + + class SubagentLifecycleService: """Stable public service returned by :attr:`PluginContext.subagent_lifecycle`. - Running children are in-process only. Completed results remain available - until process exit; ``reconnect`` accurately reports that a serialized - handle cannot reconnect after a restart instead of launching work again. + Running children are in-process only. Completed results remain available + until process exit; ``reconnect`` reports that a serialized handle cannot + reconnect after a restart instead of launching work again. """ def __init__(self, parent_agent_resolver: Callable[[], Any]) -> None: @@ -198,37 +235,26 @@ class SubagentLifecycleService: def launch(self, request: SubagentLaunchRequest) -> SubagentHandle: parent = self._parent_agent_resolver() if parent is None: - raise SubagentLifecycleError( - "No active Hermes parent session is available." - ) + raise SubagentLifecycleError("No active Hermes parent session is available.") self._validate_request(request, parent) parent_session_id = str(getattr(parent, "session_id", "") or "") or None if request.parent_session_id and request.parent_session_id != parent_session_id: - raise SubagentLifecycleError( - "parent_session_id does not match the active session." - ) + raise SubagentLifecycleError("parent_session_id does not match the active session.") correlation_key = (parent_session_id, request.correlation_id or "") with _REGISTRY.lock: self._cleanup_locked() if request.correlation_id and correlation_key in _REGISTRY.correlations: - raise SubagentLifecycleError( - "Duplicate correlation_id for this parent session." - ) + raise SubagentLifecycleError("Duplicate correlation_id for this parent session.") - # Delegate construction remains internal so plugin code never imports - # private delegation helpers or manipulates the active-child registry. - from tools.delegate_tool import ( - _build_child_preserving_parent_tools, - DEFAULT_MAX_ITERATIONS, - ) + # Delegate construction stays internal so plugin code never imports + # private delegation helpers or touches the active-child registry. + from tools.delegate_tool import _build_child_preserving_parent_tools, DEFAULT_MAX_ITERATIONS child = _build_child_preserving_parent_tools( task_index=0, goal=request.goal, context=request.context, - toolsets=list(request.allowed_toolsets) - if request.allowed_toolsets - else None, + toolsets=list(request.allowed_toolsets) if request.allowed_toolsets else None, model=request.model, max_iterations=DEFAULT_MAX_ITERATIONS, task_count=1, @@ -262,20 +288,14 @@ class SubagentLifecycleService: def status(self, handle: SubagentHandle) -> SubagentStatus: record = self._record(handle) if record is None: - return SubagentStatus( - handle, SubagentState.UNKNOWN, time.time(), "UNKNOWN_HANDLE" - ) + return SubagentStatus(handle, SubagentState.UNKNOWN, time.time(), "UNKNOWN_HANDLE") with _REGISTRY.lock: return SubagentStatus(record.handle, record.state, record.updated_at) - def wait( - self, handle: SubagentHandle, *, timeout_seconds: Optional[float] = None - ) -> SubagentTerminalState: + def wait(self, handle: SubagentHandle, *, timeout_seconds: Optional[float] = None) -> SubagentTerminalState: record = self._record(handle) if record is None: - return SubagentTerminalState( - handle, SubagentState.UNKNOWN, True, diagnostic="UNKNOWN_HANDLE" - ) + return SubagentTerminalState(handle, SubagentState.UNKNOWN, True, diagnostic="UNKNOWN_HANDLE") future = record.future if future is not None: try: @@ -285,9 +305,7 @@ class SubagentLifecycleService: except Exception: pass with _REGISTRY.lock: - return SubagentTerminalState( - record.handle, record.state, record.result is not None - ) + return SubagentTerminalState(record.handle, record.state, record.result is not None) def cancel(self, handle: SubagentHandle, *, reason: str) -> SubagentCancelResult: record = self._record(handle) @@ -295,16 +313,13 @@ class SubagentLifecycleService: return SubagentCancelResult(False, unknown_handle=True) with _REGISTRY.lock: if record.result is not None: - return SubagentCancelResult( - False, already_terminal=True, state=record.state - ) + return SubagentCancelResult(False, already_terminal=True, state=record.state) agent = record.agent record.state = SubagentState.CANCEL_REQUESTED record.updated_at = time.time() + unsupported = SubagentCancelResult(False, unsupported=True, state=SubagentState.CANCEL_REQUESTED) if agent is None: - return SubagentCancelResult( - False, unsupported=True, state=SubagentState.CANCEL_REQUESTED - ) + return unsupported try: accepted = request_hard_interrupt( agent, @@ -312,74 +327,32 @@ class SubagentLifecycleService: tool_reason="subagent cancellation requested", ) except Exception: - return SubagentCancelResult( - False, unsupported=True, state=SubagentState.CANCEL_REQUESTED - ) + accepted = False if not accepted: - return SubagentCancelResult( - False, unsupported=True, state=SubagentState.CANCEL_REQUESTED - ) + return unsupported return SubagentCancelResult(True, state=SubagentState.CANCEL_REQUESTED) def result(self, handle: SubagentHandle) -> SubagentResult: record = self._record(handle) if record is None: - return SubagentResult( - handle, - SubagentState.UNKNOWN, - False, - error_classification="UNKNOWN_HANDLE", - ) + return SubagentResult(handle, SubagentState.UNKNOWN, False, error_classification="UNKNOWN_HANDLE") with _REGISTRY.lock: if record.result is not None: return record.result - return SubagentResult( - record.handle, record.state, False, error_classification="NOT_READY" - ) + return SubagentResult(record.handle, record.state, False, error_classification="NOT_READY") def reconnect(self, handle: SubagentHandle) -> SubagentReconnectResult: record = self._record(handle) if record is None: - return SubagentReconnectResult( - False, SubagentState.UNKNOWN, "RECONNECT_UNAVAILABLE" - ) + return SubagentReconnectResult(False, SubagentState.UNKNOWN, "RECONNECT_UNAVAILABLE") with _REGISTRY.lock: return SubagentReconnectResult(True, record.state) def _record(self, handle: SubagentHandle) -> Optional[_Record]: - if ( - not isinstance(handle, SubagentHandle) - or type(handle.contract_version) is not int - or handle.contract_version != PUBLIC_CONTRACT_VERSION - ): + if not _handle_is_well_formed(handle): return None - if ( - not isinstance(handle.subagent_id, str) - or not handle.subagent_id - or ( - handle.parent_session_id is not None - and not isinstance(handle.parent_session_id, str) - ) - or ( - handle.correlation_id is not None - and not isinstance(handle.correlation_id, str) - ) - or isinstance(handle.created_at, bool) - or not isinstance(handle.created_at, (int, float)) - or not math.isfinite(handle.created_at) - or (handle.provider is not None and not isinstance(handle.provider, str)) - or (handle.model is not None and not isinstance(handle.model, str)) - or not isinstance(handle.role, str) - or type(handle.depth) is not int - or not isinstance(handle.capability, str) - ): - return None - if not hmac.compare_digest( - handle.capability, - self._capability( - handle.subagent_id, handle.parent_session_id, handle.created_at - ), - ): + expected = self._capability(handle.subagent_id, handle.parent_session_id, handle.created_at) + if not hmac.compare_digest(handle.capability, expected): return None parent = self._parent_agent_resolver() active_parent_id = str(getattr(parent, "session_id", "") or "") or None @@ -403,8 +376,7 @@ class SubagentLifecycleService: record = _REGISTRY.records.pop(subagent_id) if record.handle.correlation_id: _REGISTRY.correlations.pop( - (record.handle.parent_session_id, record.handle.correlation_id), - None, + (record.handle.parent_session_id, record.handle.correlation_id), None ) def _run(self, record: _Record, goal: str, parent: Any) -> None: @@ -417,60 +389,34 @@ class SubagentLifecycleService: from tools.delegate_tool import _run_child_lifecycle raw = _run_child_lifecycle(0, goal, record.agent, parent) - status = ( - str(raw.get("status", "error")) if isinstance(raw, dict) else "error" - ) - if status == "completed": - state = SubagentState.SUCCEEDED - elif status == "interrupted": - state = ( - SubagentState.CANCELLED - if record.state == SubagentState.CANCEL_REQUESTED - else SubagentState.INTERRUPTED - ) + is_dict = isinstance(raw, dict) + if not is_dict: + raw = {} + status = str(raw.get("status", "error")) + if status == "interrupted": + cancelled = record.state == SubagentState.CANCEL_REQUESTED + state = SubagentState.CANCELLED if cancelled else SubagentState.INTERRUPTED else: - state = SubagentState.FAILED - summary = raw.get("summary") if isinstance(raw, dict) else None - summary = str(summary)[:_MAX_RESULT_CHARS] if summary is not None else None - error = raw.get("error") if isinstance(raw, dict) else None - result = SubagentResult( - record.handle, - state, - True, - summary=summary, - completed_at=time.time(), - started_at=record.started_at, - error_classification=None - if state == SubagentState.SUCCEEDED - else status.upper(), + state = SubagentState.SUCCEEDED if status == "completed" else SubagentState.FAILED + summary = raw.get("summary") + error = raw.get("error") + fields: dict[str, Any] = dict( + summary=str(summary)[:_MAX_RESULT_CHARS] if summary is not None else None, + error_classification=None if state == SubagentState.SUCCEEDED else status.upper(), error_message=str(error)[:_MAX_RESULT_CHARS] if error else None, - usage_metadata={"api_calls": raw.get("api_calls", 0)} - if isinstance(raw, dict) - else {}, - tool_execution_summary={ - "duration_seconds": raw.get("duration_seconds", 0) - } - if isinstance(raw, dict) - else {}, + usage_metadata={"api_calls": raw.get("api_calls", 0)} if is_dict else {}, + tool_execution_summary={"duration_seconds": raw.get("duration_seconds", 0)} if is_dict else {}, ) except Exception as exc: - result = SubagentResult( - record.handle, - SubagentState.FAILED, - True, - started_at=record.started_at, - completed_at=time.time(), - error_classification=type(exc).__name__, - error_message=str(exc)[:_MAX_RESULT_CHARS], - ) + state = SubagentState.FAILED + fields = dict(error_classification=type(exc).__name__, error_message=str(exc)[:_MAX_RESULT_CHARS]) + result = SubagentResult( + record.handle, state, True, started_at=record.started_at, completed_at=time.time(), **fields + ) payload = dataclasses.asdict(result) payload.pop("result_hash", None) - result = dataclasses.replace( - result, - result_hash=hashlib.sha256( - json.dumps(payload, sort_keys=True, default=str).encode() - ).hexdigest(), - ) + digest = hashlib.sha256(json.dumps(payload, sort_keys=True, default=str).encode()).hexdigest() + result = dataclasses.replace(result, result_hash=digest) with _REGISTRY.lock: record.agent = None record.result = result @@ -479,9 +425,7 @@ class SubagentLifecycleService: record.updated_at = result.completed_at or time.time() @staticmethod - def _capability( - subagent_id: str, parent_session_id: Optional[str], created_at: float - ) -> str: + def _capability(subagent_id: str, parent_session_id: Optional[str], created_at: float) -> str: value = f"{subagent_id}|{parent_session_id or ''}|{created_at:.6f}".encode() return hmac.new(_SECRET, value, hashlib.sha256).hexdigest() @@ -493,34 +437,18 @@ class SubagentLifecycleService: or not request.goal.strip() or len(request.goal) > _MAX_GOAL_CHARS ): - raise SubagentLifecycleError( - "goal must be a non-empty string of at most 16000 characters." - ) + raise SubagentLifecycleError("goal must be a non-empty string of at most 16000 characters.") if request.context is not None and ( - not isinstance(request.context, str) - or len(request.context) > _MAX_CONTEXT_CHARS + not isinstance(request.context, str) or len(request.context) > _MAX_CONTEXT_CHARS ): - raise SubagentLifecycleError( - "context must be a string of at most 32000 characters." - ) + raise SubagentLifecycleError("context must be a string of at most 32000 characters.") if request.role not in {"leaf", "orchestrator"}: raise SubagentLifecycleError("role must be 'leaf' or 'orchestrator'.") - if request.timeout_seconds is not None: - raise SubagentLifecycleError( - "Per-launch timeout is not supported; configure delegation timeout explicitly." - ) - if request.working_directory is not None: - raise SubagentLifecycleError( - "working_directory is not supported because Hermes delegates use isolated task environments." - ) - if request.blocked_tools: - raise SubagentLifecycleError( - "Per-tool blocking is not supported; use allowed_toolsets. Hermes always blocks unsafe child tools." - ) + for rejected, message in _UNSUPPORTED_REQUEST_FIELDS: + if rejected(request): + raise SubagentLifecycleError(message) try: - metadata_bytes = len( - json.dumps(dict(request.metadata), sort_keys=True).encode() - ) + metadata_bytes = len(json.dumps(dict(request.metadata), sort_keys=True).encode()) except (TypeError, ValueError) as exc: raise SubagentLifecycleError("metadata must be JSON-serializable.") from exc if metadata_bytes > _MAX_METADATA_BYTES: @@ -530,13 +458,7 @@ class SubagentLifecycleService: unknown = set(request.allowed_toolsets) - set(TOOLSETS) if unknown: - raise SubagentLifecycleError( - f"Unknown toolsets: {', '.join(sorted(unknown))}." - ) + raise SubagentLifecycleError(f"Unknown toolsets: {', '.join(sorted(unknown))}.") enabled = getattr(parent, "enabled_toolsets", None) - if enabled is not None and not set(request.allowed_toolsets).issubset( - set(enabled) - ): - raise SubagentLifecycleError( - "Requested toolsets would broaden parent permissions." - ) + if enabled is not None and not set(request.allowed_toolsets).issubset(set(enabled)): + raise SubagentLifecycleError("Requested toolsets would broaden parent permissions.") diff --git a/agent/thread_scoped_output.py b/agent/thread_scoped_output.py index 3c4a7be891..bbbb24f3ae 100644 --- a/agent/thread_scoped_output.py +++ b/agent/thread_scoped_output.py @@ -1,19 +1,11 @@ """Thread-scoped stdout/stderr silencing for background worker threads. -``contextlib.redirect_stdout``/``redirect_stderr`` reassign the *process-global* -``sys.stdout``/``sys.stderr``. When a daemon worker thread (e.g. the background -memory/skill review) wraps its whole body in those context managers, every other -thread in the process — including a gateway's asyncio event-loop thread driving a -Telegram long-poll — sees ``sys.stdout``/``sys.stderr`` pointing at ``devnull`` -for the full duration. Any bare ``print`` / ``sys.stderr.write`` from those other -threads is silently lost during that window (see issue #55769 / #55925). - -This module installs a thin proxy as ``sys.stdout``/``sys.stderr`` that routes -writes per-thread: threads registered as "silenced" go to a sink; every other -thread passes through to the *original* stream. The proxy is installed once, -idempotently, and is never uninstalled (uninstalling would race other threads -mid-write), so the only observable effect for unregistered threads is one extra -attribute lookup per write. +``contextlib.redirect_stdout`` reassigns the *process-global* stream, so a +daemon worker silencing itself also silences every other thread (gateway +event loop included) for the duration. This module installs a per-thread +routing proxy as ``sys.stdout``/``sys.stderr``: silenced threads write to a +sink, everyone else passes through to the original stream. Installed once, +idempotently, and never uninstalled (that would race other threads mid-write). """ from __future__ import annotations @@ -27,12 +19,10 @@ from typing import Iterator, TextIO __all__ = ["thread_scoped_silence"] _install_lock = threading.Lock() -# Maps the proxy we installed for a given attribute ("stdout"/"stderr") so we -# never double-wrap and so we can recover the original stream. +# Proxy installed per attribute ("stdout"/"stderr"): never double-wrap. _installed: dict[str, "_ThreadRoutingStream"] = {} -# One process-lifetime sink per stream. Temporary process-global redirects can -# displace and later restore a routing proxy; they must not allocate another -# permanent /dev/null descriptor every time that happens. +# One process-lifetime sink per stream: global redirects that displace and +# restore a proxy must not leak a new /dev/null descriptor each time. _sinks: dict[str, TextIO] = {} _routing_states: dict[str, "_RoutingState"] = {} @@ -47,14 +37,8 @@ class _RoutingState: class _ThreadRoutingStream: - """A ``sys.stdout``/``sys.stderr`` stand-in that routes writes per-thread. - - Threads whose ident is in ``_silenced`` write to ``_sink``; all other - threads write to ``_passthrough`` (the original stream captured at install - time). Attribute access for anything other than the methods we override - is delegated to the *current* target so things like ``.encoding`` / - ``.fileno()`` behave like the underlying stream for the calling thread. - """ + """``sys.stdout``/``sys.stderr`` stand-in routing writes per calling thread; + unknown attributes delegate to the current thread's target.""" def __init__(self, passthrough: TextIO, state: _RoutingState) -> None: self._passthrough = passthrough @@ -65,7 +49,6 @@ class _ThreadRoutingStream: return self._state.sink return self._passthrough - # --- registration ----------------------------------------------------- def silence(self, ident: int) -> None: with self._state.lock: self._state.silenced[ident] = self._state.silenced.get(ident, 0) + 1 @@ -78,7 +61,6 @@ class _ThreadRoutingStream: else: self._state.silenced.pop(ident, None) - # --- file-like surface ------------------------------------------------ def write(self, data): # type: ignore[no-untyped-def] try: return self._target().write(data) @@ -108,8 +90,6 @@ class _ThreadRoutingStream: return self._target().fileno() def __getattr__(self, name): # type: ignore[no-untyped-def] - # Delegate everything we don't override (encoding, buffer, mode, ...) - # to the calling thread's current target. return getattr(self._target(), name) @@ -119,17 +99,15 @@ def _ensure_installed(attr: str, passthrough: TextIO) -> "_ThreadRoutingStream": proxy = _installed.get(attr) current = getattr(sys, attr, None) if isinstance(current, _ThreadRoutingStream): - # A redirect context can restore an older routing proxy after a - # temporary replacement. Adopt it instead of wrapping it and - # growing an unbounded proxy chain. + # A redirect context may restore an older proxy; adopt it rather + # than wrapping it into an unbounded chain. _installed[attr] = current _routing_states[attr] = current._state return current if proxy is not None and current is proxy: return proxy - # Capture whatever is currently bound as the passthrough. If a prior - # global redirect_stdout is active, route non-silenced threads to that - # stream to preserve the old behavior. + # Route non-silenced threads to whatever is currently bound (an active + # global redirect keeps its old behavior). passthrough = current if current is not None else passthrough sink = _sinks.get(attr) if sink is None or sink.closed: @@ -147,12 +125,7 @@ def _ensure_installed(attr: str, passthrough: TextIO) -> "_ThreadRoutingStream": @contextlib.contextmanager def thread_scoped_silence() -> Iterator[None]: - """Silence ``stdout``/``stderr`` for the *current thread only*. - - Other threads keep writing to the real streams. Use this around a worker - thread's body instead of ``contextlib.redirect_stdout(devnull)`` when the - process is multi-threaded and another thread must keep its console output. - """ + """Silence ``stdout``/``stderr`` for the *current thread only*.""" ident = threading.get_ident() out_proxy = _ensure_installed("stdout", sys.__stdout__ or sys.stdout) err_proxy = _ensure_installed("stderr", sys.__stderr__ or sys.stderr) diff --git a/agent/turn_liveness.py b/agent/turn_liveness.py index 939fc9b917..04e381df1f 100644 --- a/agent/turn_liveness.py +++ b/agent/turn_liveness.py @@ -1,57 +1,20 @@ -"""Turn liveness watchdog (#95548): force-abort turns that stall silently. +"""Turn liveness watchdog: force-abort turns that stall silently. -A conversation turn can stall mid-flight (observed in #95548 between -"model returned tool_calls" and tool execution, after a slow model -response + desktop WS disconnect) with no error logged, no further -progress, and the durable session turn lease kept renewing — so nothing -ever force-aborts the turn and the session stays stuck until the process -is killed. +A turn can wedge mid-flight with no error and its durable lease still +renewing, so nothing ever frees the session. This module owns the policy: +config resolution (``agent.turn_liveness`` in config.yaml, validated — a typo, +NaN or Inf warns and falls back rather than disabling the timeout or freezing +the poll), the sampling state machine, and the watcher thread. +``AIAgent.run_conversation`` supplies the commit/deactivate callbacks that own +turn-lease state. -This module owns the watchdog policy end to end: - -* configuration — :func:`resolve_turn_liveness_settings` reads the - ``agent.turn_liveness`` section of config.yaml and validates it; -* state machine — :class:`TurnLivenessWatchdog` samples the agent's - activity clock and drives the stall decision; -* thread mechanics — the polling loop, stop handling, and the stall - commit point. - -``AIAgent.run_conversation`` (``run_agent.py``) keeps only the smallest -integration seam: it resolves the settings, supplies the commit and -deactivate callbacks that own turn-lease state, and starts the thread -next to the durable lease refresher. - -Config surface (config.yaml):: - - agent: - turn_liveness: - timeout_s: 600.0 # idle bound; <= 0 disables the watchdog - poll_s: 15.0 # sampling interval (seconds) - -Both values are validated. A non-numeric typo, ``NaN``, or ``Inf`` logs a -warning and falls back to the documented default — it never crashes -startup, and a bogus value can never silently disable the watchdog -(``NaN``) or freeze the watcher thread (``Inf`` poll). - -Race safety (#95663 review): the watchdog samples the activity clock and -binds the abort decision to the observed ``(generation, timestamp)`` pair. -The commit callback revalidates that pair under the *same* lock -``AIAgent._touch_activity`` uses to stamp the clock, so a turn that -resumed while the stall was being logged/emitted is never hard-cancelled -— it continues and its lease keeps renewing. The revalidated generation -is carried into ``AIAgent.interrupt`` (``require_generation``) as a -cancellation claim consumed at the final mutation edge: ``interrupt`` -reserves the claim under the activity lock, ``_touch_activity`` -invalidates the reservation the instant real progress lands, and the -claim survives every blocking boundary, including the compression -commit fence. Claim consumption and the first observable interrupt -state publish in ONE activity-lock critical section, so a turn that -resumes can only ever interleave before that section (the reservation -is invalidated and the abort declines) or after it (the interrupt -already committed under the lock) — never between "claim consumed" and -"state published". A turn that resumes anywhere in the window is never -hard-cancelled, and an exceptional interrupt path declines the abort -fail-closed instead of mutating interrupt state. +Race safety: the watchdog binds its abort decision to the observed +``(generation, timestamp)`` pair; the commit callback revalidates that pair +under the same lock ``_touch_activity`` stamps with, so a turn that resumed +while the stall was being surfaced is never hard-cancelled. The revalidated +generation flows into ``AIAgent.interrupt`` (``require_generation``) as a +claim consumed at the final mutation edge, in ONE activity-lock critical +section with the first interrupt-state publish. """ from __future__ import annotations @@ -73,13 +36,8 @@ _CONFIG_POLL_KEY = "agent.turn_liveness.poll_s" class ActivitySnapshot(NamedTuple): - """One observation of the activity clock, bound to an abort decision. - - ``generation`` + ``activity_ts`` uniquely identify the observed stamp. - The commit callback must revalidate this pair under the lock shared - with ``AIAgent._touch_activity``; if it no longer matches, the - observation is stale and the abort must be declined. - """ + """One activity-clock observation; ``(generation, activity_ts)`` must be + revalidated by the commit callback under the shared lock.""" generation: int activity_ts: Optional[float] @@ -87,27 +45,19 @@ class ActivitySnapshot(NamedTuple): def _warn_invalid_value(key: str, raw: Any, default: float) -> None: - logger.warning( - "Invalid %s in config.yaml: %r — falling back to default %.1f.", - key, - raw, - default, - ) + logger.warning("Invalid %s in config.yaml: %r — falling back to default %.1f.", key, raw, default) def _resolve_finite_seconds(raw: Any, *, default: float, key: str) -> float: - """Coerce one duration knob, rejecting typos, NaN and Inf. + """Coerce one duration knob, rejecting typos, NaN and Inf (never raises). - A non-numeric value must not raise into durable-turn startup, and a - non-finite value must not silently change behavior (``NaN`` would - disable the timeout via the ``> 0`` comparison; ``Inf`` would freeze - the poll loop in ``Event.wait``). + NaN would silently disable the timeout via the ``> 0`` comparison and + Inf would freeze the poll loop in ``Event.wait``. """ try: value = float(raw) except (TypeError, ValueError): - _warn_invalid_value(key, raw, default) - return default + value = math.nan if not math.isfinite(value): _warn_invalid_value(key, raw, default) return default @@ -117,25 +67,16 @@ def _resolve_finite_seconds(raw: Any, *, default: float, key: str) -> float: def resolve_turn_liveness_settings( config: Optional[Dict[str, Any]] = None, ) -> Tuple[Optional[float], float]: - """Resolve ``(timeout_s, poll_s)`` from the ``agent.turn_liveness`` section. + """Resolve ``(timeout_s, poll_s)``; ``timeout_s <= 0`` opts out (``None``). - Precedence: explicit config.yaml value wins over the default; any - invalid value (typo, NaN, Inf, non-positive poll) falls back to the - default with a warning. ``timeout_s <= 0`` is the documented opt-out - and yields ``(None, poll_s)`` — the caller then never arms the - watchdog. The resolver never raises. + Invalid values (typo, NaN, Inf, non-positive poll) warn and fall back to + the default. Never raises. """ - section: Dict[str, Any] = {} - if isinstance(config, dict): - agent_cfg = config.get("agent") - if isinstance(agent_cfg, dict): - raw_section = agent_cfg.get("turn_liveness") - if isinstance(raw_section, dict): - section = raw_section - elif raw_section is not None: - _warn_invalid_value( - "agent.turn_liveness", raw_section, DEFAULT_TURN_LIVENESS_TIMEOUT_S - ) + agent_cfg = config.get("agent") if isinstance(config, dict) else None + raw_section = agent_cfg.get("turn_liveness") if isinstance(agent_cfg, dict) else None + section: Dict[str, Any] = raw_section if isinstance(raw_section, dict) else {} + if raw_section is not None and not isinstance(raw_section, dict): + _warn_invalid_value("agent.turn_liveness", raw_section, DEFAULT_TURN_LIVENESS_TIMEOUT_S) timeout_s = _resolve_finite_seconds( section.get("timeout_s", DEFAULT_TURN_LIVENESS_TIMEOUT_S), @@ -150,7 +91,6 @@ def resolve_turn_liveness_settings( if poll_s <= 0: _warn_invalid_value(_CONFIG_POLL_KEY, poll_s, DEFAULT_TURN_LIVENESS_POLL_S) poll_s = DEFAULT_TURN_LIVENESS_POLL_S - # <= 0 is the documented opt-out, not an error. if timeout_s <= 0: timeout_s = None return timeout_s, poll_s @@ -159,11 +99,8 @@ def resolve_turn_liveness_settings( class TurnLivenessWatchdog: """Sampled-idle watchdog thread bound to one conversation turn. - ``run_agent.py`` owns the turn-lease state (stop event, turn-active - flag, interrupt plumbing); this class only reads the activity clock - and calls back at the stall commit point. All synchronization with - ``AIAgent._touch_activity`` goes through ``activity_lock``, which - must be the SAME lock the agent stamps its activity clock with. + ``activity_lock`` must be the SAME lock ``AIAgent._touch_activity`` stamps + the activity clock with; run_agent owns the lease state and callbacks. """ def __init__( @@ -190,55 +127,26 @@ class TurnLivenessWatchdog: self._deactivate_turn = deactivate_turn def make_thread(self) -> threading.Thread: - """Build the (not yet started) watcher thread. - - ``run_agent.py`` creates the watchdog before the turn begins but - starts the thread at turn entry, right after the turn-active flag - and the activity clock are stamped. - """ - return threading.Thread( - target=self._watch, - name="turn-liveness-watchdog", - daemon=True, - ) - - def start(self) -> threading.Thread: - """Spawn the watcher thread and return it (already running).""" - thread = self.make_thread() - thread.start() - return thread + """Build the (not yet started) watcher thread; run_agent starts it at + turn entry, after the turn-active flag and activity clock are stamped.""" + return threading.Thread(target=self._watch, name="turn-liveness-watchdog", daemon=True) def _watch(self) -> None: while not self._stop_event.wait(self._poll_s): snapshot = self._sample() if snapshot is None: - # Turn is no longer active; nothing to watch. - return + return # turn no longer active if snapshot.idle_seconds < self._timeout_s: continue - # Pre-commit surface is OBSERVATIONAL only: it reports the - # stall and that a recovery attempt is beginning. It must not - # claim the abort or the lease withdrawal has committed — the - # next operation can still veto the outcome. The definitive - # aborted/lease-stopped settlement is published by - # _surface_committed_abort only after _commit_abort succeeds - # and the turn is deactivated (#95663 review). + # Observational only: the commit below can still veto the abort + # if progress resumed; the definitive settlement is published + # by _surface_committed_abort after commit + deactivate. self._surface_stall(snapshot) - # Commit point: bind the abort to the sampled generation/ts - # and revalidate under the lock shared with `_touch_activity`. - # If progress resumed while the stall was being surfaced, the - # turn continues and this loop resumes sampling — the lease - # keeps renewing. The commit also carries the revalidated - # generation into the interrupt path, which reserves it as a - # claim, survives every blocking boundary (compression - # fence), and consumes it at the final mutation edge — progress - # landing anywhere in that window declines the abort. - if not self._commit_abort(snapshot, self._abort_message(snapshot)): + message = f"Turn made no progress for {int(snapshot.idle_seconds)}s; aborting to release the session." + if not self._commit_abort(snapshot, message): continue - # Stop renewing the durable lease: a wedge the hard interrupt - # cannot unwind must not keep the lease alive forever (the - # issue's "lease keeps renewing" masking). The TTL expiry then - # lets stale-turn cleanup reclaim the row. + # Stop renewing the lease so a wedge the interrupt cannot unwind + # expires via TTL instead of masking forever. self._deactivate_turn() self._surface_committed_abort(snapshot) return @@ -247,97 +155,59 @@ class TurnLivenessWatchdog: with self._activity_lock: if not self._is_turn_active(): return None - generation = getattr( - self._agent, "_turn_liveness_activity_generation", 0 - ) + generation = getattr(self._agent, "_turn_liveness_activity_generation", 0) activity_ts = getattr(self._agent, "_last_activity_ts", None) - now = time.time() - if activity_ts is None: - idle_seconds = 0.0 - else: - idle_seconds = max(0.0, now - activity_ts) - return ActivitySnapshot( - generation=generation, - activity_ts=activity_ts, - idle_seconds=idle_seconds, - ) + idle_seconds = 0.0 if activity_ts is None else max(0.0, time.time() - activity_ts) + return ActivitySnapshot(generation, activity_ts, idle_seconds) - def _abort_message(self, snapshot: ActivitySnapshot) -> str: - return ( - f"Turn made no progress for {int(snapshot.idle_seconds)}s; " - "aborting to release the session." - ) + def _emit_warning(self, text: str, debug_msg: str) -> None: + emit_warning = getattr(self._agent, "_emit_warning", None) + if not callable(emit_warning): + return + try: + emit_warning(text) + except Exception: + logger.debug(debug_msg, exc_info=True) def _surface_stall(self, snapshot: ActivitySnapshot) -> None: - """Observationally surface the stall: log it loudly and emit a - UI-visible warning that a recovery attempt is beginning. + """Log + UI-warn that recovery is beginning (not that it committed). - Deliberately does NOT claim the abort or the lease withdrawal has - committed: the next operation (``_commit_abort``) can still veto - the outcome when the turn resumed while this surface window was - open. The definitive aborted/lease-stopped settlement is - published by :meth:`_surface_committed_abort` only after the - abort wins and the turn is deactivated. - - Rate-limited: a turn whose aborts keep declining (resumed - activity, exceptional interrupt path) must not re-log and - re-warn every poll interval — the first surface carries the - signal, repeats are suppressed until activity actually moves - again (a new generation re-arms the surface). + Rate-limited per activity generation so repeatedly declined aborts + do not re-log every poll; a new generation re-arms the surface. """ generation = snapshot.generation if getattr(self, "_last_surfaced_generation", None) == generation: return self._last_surfaced_generation = generation - session_id = getattr(self._agent, "session_id", None) or self._session_id - last_desc = getattr(self._agent, "_last_activity_desc", None) logger.error( "Turn liveness watchdog fired for session %s: " "no progress for %.1fs (last activity: %r). " "Attempting recovery: force-interrupting the turn and " "stopping lease renewal if it cannot resume (#95548).", - session_id, + getattr(self._agent, "session_id", None) or self._session_id, snapshot.idle_seconds, - last_desc, + getattr(self._agent, "_last_activity_desc", None), + ) + self._emit_warning( + "⚠️ This turn stopped making progress " + f"({int(snapshot.idle_seconds)}s without activity); " + "attempting recovery so the session can continue.", + "Failed to emit turn liveness warning", ) - emit_warning = getattr(self._agent, "_emit_warning", None) - if not callable(emit_warning): - return - try: - emit_warning( - "⚠️ This turn stopped making progress " - f"({int(snapshot.idle_seconds)}s without activity); " - "attempting recovery so the session can continue." - ) - except Exception: - logger.debug("Failed to emit turn liveness warning", exc_info=True) def _surface_committed_abort(self, snapshot: ActivitySnapshot) -> None: - """Publish the definitive settlement AFTER the abort has authority. - - Runs only once ``_commit_abort`` succeeded (the interrupt was - published) and the turn lease was deactivated: the turn IS - force-aborted and lease renewal IS stopped, so stating that is - now true. Separated from the pre-commit surface so a declined - abort never reports a committed outcome (#95663 review). - """ - session_id = getattr(self._agent, "session_id", None) or self._session_id + """Publish the definitive settlement once the abort has authority.""" logger.error( "Turn liveness watchdog aborted turn for session %s: " "no progress for %.1fs; turn interrupted and lease renewal " "stopped (#95548).", - session_id, + getattr(self._agent, "session_id", None) or self._session_id, snapshot.idle_seconds, ) - emit_warning = getattr(self._agent, "_emit_warning", None) - if not callable(emit_warning): - return - try: - emit_warning( - "⚠️ Turn aborted by the liveness watchdog " - f"({int(snapshot.idle_seconds)}s without activity); " - "lease renewal stopped so the session can be reclaimed. " - "You can retry your message." - ) - except Exception: - logger.debug("Failed to emit committed-abort warning", exc_info=True) + self._emit_warning( + "⚠️ Turn aborted by the liveness watchdog " + f"({int(snapshot.idle_seconds)}s without activity); " + "lease renewal stopped so the session can be reclaimed. " + "You can retry your message.", + "Failed to emit committed-abort warning", + ) diff --git a/tests/agent/test_deadline.py b/tests/agent/test_deadline.py index 5122b0f26a..97b288d969 100644 --- a/tests/agent/test_deadline.py +++ b/tests/agent/test_deadline.py @@ -170,7 +170,6 @@ class TestRunBoundedSync: result = run_bounded_sync(lambda: "ok", 5.0, label="t") assert result.timed_out is False assert result.value == "ok" - assert result.raise_if_timed_out() == "ok" def test_unbounded_when_timeout_none(self): result = run_bounded_sync(lambda: 7, None, label="t") @@ -196,9 +195,7 @@ class TestRunBoundedSync: assert result.timed_out is True assert result.value is None assert elapsed < 5.0 # returned near the deadline, not after 30s - with pytest.raises(DeadlineExpired) as exc_info: - result.raise_if_timed_out() - assert "wedged" in str(exc_info.value) + assert result.label == "wedged" release.set() def test_on_timeout_callback_runs(self): diff --git a/tests/agent/test_file_safety_cross_profile.py b/tests/agent/test_file_safety_cross_profile.py index d3a68b6fb5..3be8ea8791 100644 --- a/tests/agent/test_file_safety_cross_profile.py +++ b/tests/agent/test_file_safety_cross_profile.py @@ -1,15 +1,4 @@ -"""Tests for the cross-Hermes-profile write guard in agent/file_safety. - -The guard fires when a tool tries to write into another Hermes profile's -skills/plugins/cron/memories directory. It's a soft guard — defense in -depth, NOT a security boundary — but it prevents the agent from silently -corrupting a profile that belongs to a different session. - -Reference: May 2026 incident — a hermes-security profile session -accidentally edited skills under both ~/.hermes/profiles/hermes-security/skills/ -AND ~/.hermes/skills/ (the default profile's skills), realizing only -afterwards that the second path belonged to a different profile. -""" +"""Tests for the (retired) cross-Hermes-profile write guard in agent/file_safety.""" from __future__ import annotations from pathlib import Path @@ -60,7 +49,6 @@ def fake_hermes(tmp_path, monkeypatch): import hermes_constants monkeypatch.setattr(hermes_constants, "get_default_hermes_root", lambda: root) - # The reloads below ensure get_cross_profile_warning/classify see the patched root. import agent.file_safety as fs monkeypatch.setattr(fs, "_hermes_root_path", lambda: root) @@ -102,49 +90,6 @@ class TestResolveActiveProfileName: assert fs._resolve_active_profile_name() == "default" -# --------------------------------------------------------------------------- -# classify_cross_profile_target -# --------------------------------------------------------------------------- - - -class TestClassifyCrossProfileTarget: - - def test_security_writing_default_skill(self, fake_hermes, monkeypatch): - """The exact incident from May 2026.""" - _set_active_home(monkeypatch, fake_hermes["security_home"]) - from agent.file_safety import classify_cross_profile_target - result = classify_cross_profile_target( - str(fake_hermes["default_home"] / "skills" / "foo" / "SKILL.md") - ) - assert result is not None - assert result["active_profile"] == "hermes-security" - assert result["target_profile"] == "default" - assert result["area"] == "skills" - - def test_default_writing_security_skill(self, fake_hermes, monkeypatch): - """Inverse direction — default-profile session reaching into a named profile.""" - _set_active_home(monkeypatch, fake_hermes["default_home"]) - from agent.file_safety import classify_cross_profile_target - result = classify_cross_profile_target( - str(fake_hermes["security_home"] / "skills" / "foo" / "SKILL.md") - ) - assert result is not None - assert result["active_profile"] == "default" - assert result["target_profile"] == "hermes-security" - - - @pytest.mark.parametrize("area", ["skills", "plugins", "cron", "memories"]) - def test_all_profile_scoped_areas_classified(self, fake_hermes, monkeypatch, area): - _set_active_home(monkeypatch, fake_hermes["security_home"]) - from agent.file_safety import classify_cross_profile_target - target = fake_hermes["default_home"] / area / "foo.txt" - result = classify_cross_profile_target(str(target)) - assert result is not None - assert result["area"] == area - - - - # --------------------------------------------------------------------------- # get_cross_profile_warning # --------------------------------------------------------------------------- diff --git a/tests/agent/test_file_safety_sandbox_mirror.py b/tests/agent/test_file_safety_sandbox_mirror.py index bb59c1ecfb..27e120088f 100644 --- a/tests/agent/test_file_safety_sandbox_mirror.py +++ b/tests/agent/test_file_safety_sandbox_mirror.py @@ -155,9 +155,5 @@ class TestSandboxMirrorIsOrthogonalToCrossProfile: target.parent.mkdir(parents=True) target.write_text("x") - # cross-profile classifier: active profile == target's inner-mirror - # profile name; on the existing detector the path's parts[2] is - # ``sandboxes``, not a scoped area, so it returns None. - assert fs.classify_cross_profile_target(str(target)) is None # sandbox-mirror classifier: fires unconditionally on the shape. assert fs.classify_sandbox_mirror_target(str(target)) is not None diff --git a/tests/test_estop.py b/tests/test_estop.py index e6917c7d6d..85abbfaaab 100644 --- a/tests/test_estop.py +++ b/tests/test_estop.py @@ -25,7 +25,7 @@ from agent import estop def hermes_home(tmp_path, monkeypatch): """Point HERMES_HOME at a temp dir and reset estop module log state.""" monkeypatch.setenv("HERMES_HOME", str(tmp_path)) - estop._reset_log_state_for_tests() + estop._logged_components.clear() return tmp_path @@ -389,7 +389,7 @@ def test_profile_gateway_honors_canonical_root_estop(tmp_path, monkeypatch): profile = root / "profiles" / "fleet-analyst" profile.mkdir(parents=True) monkeypatch.setenv("HERMES_HOME", str(profile)) - estop._reset_log_state_for_tests() + estop._logged_components.clear() assert estop.is_engaged() is False (root / "ESTOP").write_text("{\"reason\": \"thundering herd\"}\n", encoding="utf-8")