Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 84 additions & 11 deletions hermes_cli/plugins.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down Expand Up @@ -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())
Expand All @@ -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] = {}

Expand All @@ -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
Expand All @@ -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:
Expand Down
114 changes: 114 additions & 0 deletions tests/hermes_cli/test_plugins.py
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading