From 69e71df79303e85d04b10df07a7010342d0d6a04 Mon Sep 17 00:00:00 2001 From: Kyzcreig <9063726+Kyzcreig@users.noreply.github.com> Date: Mon, 21 Sep 2026 03:38:46 -0700 Subject: [PATCH 1/2] fix(plugins): isolate policy hook dispatch Run pre_tool_call callbacks on a dedicated daemon executor, measure timeout from callback start, and reserve suppression for callback runtime breaches. Log queue wait and runtime separately.\n\nVerified: scripts/run_tests.sh tests/agent/test_shell_hooks.py tests/hermes_cli/test_plugins.py -q (128 passed). --- hermes_cli/plugins.py | 95 +++++++++++++++++++++++--- tests/hermes_cli/test_plugins.py | 114 +++++++++++++++++++++++++++++++ 2 files changed, 198 insertions(+), 11 deletions(-) diff --git a/hermes_cli/plugins.py b/hermes_cli/plugins.py index 2fd875cc5db86..8102a46c9b9f1 100644 --- a/hermes_cli/plugins.py +++ b/hermes_cli/plugins.py @@ -457,6 +457,29 @@ def _install_plugin_debug_handler(force: bool = False) -> None: # repeatedly invoked hung hook cannot accumulate abandoned daemon threads. _HOOK_TIMEOUT_SUPPRESSION_SECONDS = 60.0 +# Policy hooks use an isolated daemon pool. They must not queue behind asyncio's +# process-wide default executor, which is also used by cron and tool work. Four +# workers bound damage from distinct hung policies while allowing concurrent +# sessions to run fast policy checks independently. +_HOOK_POLICY_EXECUTOR_WORKERS = 4 +_hook_policy_executor: Any = None +_hook_policy_executor_lock = threading.Lock() + + +def _get_hook_callback_executor(): + global _hook_policy_executor + if _hook_policy_executor is None: + with _hook_policy_executor_lock: + if _hook_policy_executor is None: + from tools.daemon_pool import DaemonThreadPoolExecutor + + _hook_policy_executor = DaemonThreadPoolExecutor( + max_workers=_HOOK_POLICY_EXECUTOR_WORKERS, + thread_name_prefix="hermes-policy-hook", + ) + return _hook_policy_executor + + _PRE_TOOL_CALL_TIMEOUT_BLOCK_MESSAGE = ( "pre_tool_call plugin callback timed out or is still running" ) @@ -5688,14 +5711,20 @@ def invoke_hook(self, hook_name: str, **kwargs: Any) -> List[Any]: callback_key ) running = callback_key in self._hook_running_callbacks - if ( + suppressed = ( suppressed_until is not None and suppressed_until > now - ) or running: + ) + if suppressed or running: + reason = ( + "suppressed after callback runtime timeout" + if suppressed + else "skipped while callback is still running" + ) logger.warning( - "Hook '%s' callback %s skipped after previous " - "timeout or while still running", + "Hook '%s' callback %s %s", hook_name, callback_name, + reason, ) if fail_closed: results.append(_pre_tool_call_timeout_block()) @@ -5705,7 +5734,10 @@ def invoke_hook(self, hook_name: str, **kwargs: Any) -> List[Any]: self._hook_running_callbacks[callback_key] = token context = contextvars.copy_context() + submitted_at = time.monotonic() + started = threading.Event() done = threading.Event() + timing: Dict[str, float] = {} outcome: Dict[str, Any] = {} failure: Dict[str, Exception] = {} @@ -5714,6 +5746,8 @@ def _runner( _key: tuple = callback_key, _token: object = token, ) -> None: + timing["started_at"] = time.monotonic() + started.set() try: # Route through _invoke_hook_callback so the # additive-payload signature filtering (narrow @@ -5724,28 +5758,67 @@ def _runner( except Exception as exc: failure["exc"] = exc finally: + timing["finished_at"] = time.monotonic() with self._hook_timeout_lock: if self._hook_running_callbacks.get(_key) is _token: self._hook_running_callbacks.pop(_key, None) done.set() - thread = threading.Thread( - target=_runner, - name=f"hermes-hook-{callback_name}"[:40], - daemon=True, - ) - thread.start() + future = None + if fail_closed: + # Policy callbacks need a pool isolated from unrelated + # asyncio/cron/tool work. Daemon workers preserve the + # abandon-without-join timeout contract. + future = _get_hook_callback_executor().submit(_runner) + else: + thread = threading.Thread( + target=_runner, + name=f"hermes-hook-{callback_name}"[:40], + daemon=True, + ) + thread.start() + + if not started.wait(timeout=timeout): + queue_wait = time.monotonic() - submitted_at + cancelled = future is not None and future.cancel() + if cancelled: + with self._hook_timeout_lock: + if self._hook_running_callbacks.get(callback_key) is token: + self._hook_running_callbacks.pop(callback_key, None) + logger.warning( + "Hook '%s' callback %s dispatch timed out before start: " + "queue_wait=%.3fs run_time=0s budget=%gs — skipping", + hook_name, + callback_name, + queue_wait, + timeout, + ) + # Dispatch starvation is not callback slowness. Never + # start the 60s suppression window for work that did + # not execute; a later call gets a fresh policy check. + if fail_closed: + results.append(_pre_tool_call_timeout_block()) + continue + + started_at = timing["started_at"] + queue_wait = started_at - submitted_at if not done.wait(timeout=timeout): # Do not join — that would reintroduce the #6622 hang. + # Suppression is reserved for callbacks that actually + # started and exhausted their own runtime budget. + run_time = time.monotonic() - started_at with self._hook_timeout_lock: self._hook_timeout_suppressed_until[callback_key] = ( time.monotonic() + self._hook_timeout_suppression_seconds ) logger.warning( - "Hook '%s' callback %s timed out after %gs — skipping", + "Hook '%s' callback %s timed out: queue_wait=%.3fs " + "run_time=%.3fs budget=%gs — skipping", hook_name, callback_name, + queue_wait, + run_time, timeout, ) if fail_closed: diff --git a/tests/hermes_cli/test_plugins.py b/tests/hermes_cli/test_plugins.py index 958edf8c5b30d..7a79dfca02002 100644 --- a/tests/hermes_cli/test_plugins.py +++ b/tests/hermes_cli/test_plugins.py @@ -1,5 +1,7 @@ """Tests for the Hermes plugin system (hermes_cli.plugins).""" +import asyncio +from concurrent.futures import Future, ThreadPoolExecutor import logging import json import sys @@ -1073,6 +1075,118 @@ def test_hook_callback_within_timeout_returns_value(self, monkeypatch): {"context": "hi"} ] + def test_pre_tool_call_ignores_saturated_default_executor(self, monkeypatch): + """Unrelated asyncio work must not consume the policy hook's dispatch lane.""" + monkeypatch.setattr( + "hermes_cli.plugins._resolve_hook_callback_timeout", lambda: 0.2 + ) + loop = asyncio.new_event_loop() + default_pool = ThreadPoolExecutor(max_workers=1) + hold = threading.Event() + loop.set_default_executor(default_pool) + blocker = loop.run_in_executor(None, hold.wait, 10.0) + loop.run_until_complete(asyncio.sleep(0)) + + calls = [] + callback_threads = [] + + def fast_policy(**_kwargs): + threading.Event().wait(0.03) + calls.append(1) + callback_threads.append(threading.current_thread().name) + return None + + mgr = PluginManager() + mgr._hooks["pre_tool_call"] = [fast_policy] + try: + assert mgr.invoke_hook("pre_tool_call", tool_name="terminal") == [] + assert calls == [1] + assert callback_threads[0].startswith("hermes-policy-hook") + finally: + hold.set() + loop.run_until_complete(blocker) + default_pool.shutdown(wait=True) + loop.close() + + def test_dispatch_timeout_does_not_suppress_next_policy_call( + self, monkeypatch, caplog + ): + """A callback that never started must not poison the 60s suppression map.""" + from hermes_cli import plugins as plugins_mod + from hermes_cli.plugins import _PRE_TOOL_CALL_TIMEOUT_BLOCK_MESSAGE + + monkeypatch.setattr( + plugins_mod, "_resolve_hook_callback_timeout", lambda: 0.05 + ) + + class QueuedExecutor: + def __init__(self): + self.run_inline = False + + def submit(self, callback): + future = Future() + if self.run_inline: + try: + future.set_result(callback()) + except BaseException as exc: + future.set_exception(exc) + return future + + executor = QueuedExecutor() + monkeypatch.setattr( + plugins_mod, + "_get_hook_callback_executor", + lambda: executor, + raising=False, + ) + calls = [] + + def policy(**_kwargs): + calls.append(1) + return None + + mgr = PluginManager() + mgr._hooks["pre_tool_call"] = [policy] + + with caplog.at_level(logging.WARNING, logger="hermes_cli.plugins"): + first = mgr.invoke_hook("pre_tool_call", tool_name="terminal") + + assert first == [ + {"action": "block", "message": _PRE_TOOL_CALL_TIMEOUT_BLOCK_MESSAGE} + ] + assert calls == [] + assert mgr._hook_timeout_suppressed_until == {} + assert "queue_wait=" in caplog.text + assert "run_time=0" in caplog.text + + executor.run_inline = True + assert mgr.invoke_hook("pre_tool_call", tool_name="terminal") == [] + assert calls == [1] + + def test_running_hook_timeout_logs_queue_and_runtime(self, monkeypatch, caplog): + from hermes_cli import plugins as plugins_mod + + monkeypatch.setattr( + plugins_mod, "_resolve_hook_callback_timeout", lambda: 0.05 + ) + hold = threading.Event() + started = threading.Event() + + def hung_policy(**_kwargs): + started.set() + hold.wait(timeout=10.0) + + mgr = PluginManager() + mgr._hooks["pre_tool_call"] = [hung_policy] + try: + with caplog.at_level(logging.WARNING, logger="hermes_cli.plugins"): + mgr.invoke_hook("pre_tool_call", tool_name="terminal") + assert started.is_set() + assert "queue_wait=" in caplog.text + assert "run_time=" in caplog.text + finally: + hold.set() + def test_hook_exception_still_isolated_under_timeout_path(self, monkeypatch): monkeypatch.setattr( "hermes_cli.plugins._resolve_hook_callback_timeout", lambda: 1.0 From 2462810c09a390c4c27034a1240592d80c62c439 Mon Sep 17 00:00:00 2001 From: Alexander Nikolas Date: Mon, 21 Sep 2026 04:06:35 -0700 Subject: [PATCH 2/2] chore: re-submit head for FleetReview (69e71df7 was stamped ERROR by the 10:42Z router restart, not a verdict; kanban review APPROVED at 69e71df7, tree identical)