fix(plugins): ainvoke_hook shares the sync path's failure contract

Follow-up trim of the #110265 salvage. `ainvoke_hook` logged raising callbacks
with a bare warning; route them through `_report_hook_failure` (warn-once per
distinct failure, #111922) and, for `_HOOK_TIMEOUT_FAIL_CLOSED_HOOKS`, append
the same named block directive the sync path emits (#109624), so the async twin
cannot drift into a fail-open policy path. Tests trimmed to the salvage bar: the
in-loop await is proven once through the real `_handle_message` path
(`test_async_hook_callback_is_awaited_on_the_gateway_loop`); the manager-level
duplicate is dropped and the narrowing test also pins failure isolation. Docs:
`pre_gateway_dispatch` callbacks may be `async def` and stay unbounded.

Credit order for the three PRs fixing this gap: #102485 (dmspark, earliest,
pre-decomposition `gateway/run.py`), #110253 (KoNit-K, bounded the hook —
rejected by design: neither fail mode is acceptable for a policy gate), #110265
(twidtwid, reporter; cherry-picked because it matches the ainvoke_hook shape,
keeps the hook unbounded, and adapts the existing sync test seams honestly).

Part of #110241
Supersedes #102485
Supersedes #110253
Co-authored-by: David Marcus <dmspark@users.noreply.github.com>
Co-authored-by: KoNit-K <konit.block@protonmail.com>
This commit is contained in:
teknium1
2026-09-21 22:46:43 -07:00
committed by Teknium
parent 727e7342cf
commit 75e9567ca7
3 changed files with 20 additions and 30 deletions

View File

@@ -510,8 +510,12 @@ class PluginDispatchMixin:
logger.warning("Hook '%s' callback %s timed out after %.0fs", hook_name, callback_name, timeout)
if fail_closed: # policy hook: fail closed with a block directive
results.append({"action": "block", "message": _PRE_TOOL_CALL_TIMEOUT_BLOCK_MESSAGE})
except Exception as exc:
logger.warning("Hook '%s' callback %s raised: %s", hook_name, callback_name, exc)
except (Exception, SystemExit) as exc:
# Same isolation + failure contract as the sync path (#111922 warn-once, #109624
# a raising policy guard fails closed).
self._report_hook_failure(hook_name, cb, kwargs, exc)
if fail_closed:
results.append(_policy_error_block_directive(hook_name, cb, exc))
return results
def iter_hook_callbacks(self, hook_name: str) -> tuple[Callable, ...]:

View File

@@ -2888,32 +2888,11 @@ class TestAsyncHookOnCallerLoop:
the callback on the caller's loop.
"""
def test_callback_that_needs_the_caller_loop_completes(self):
import asyncio
mgr = PluginManager()
async def driver():
gate = asyncio.Event()
async def async_hook(**kwargs):
await gate.wait() # only a sibling task on THIS loop can release it
return {"action": "allow"}
async def release():
await asyncio.sleep(0)
gate.set()
mgr._hooks.setdefault("pre_gateway_dispatch", []).append(async_hook)
asyncio.create_task(release())
return await asyncio.wait_for(
mgr.ainvoke_hook("pre_gateway_dispatch", event="e", gateway="g"), timeout=5)
assert asyncio.run(driver()) == [{"action": "allow"}]
def test_narrow_legacy_signature_still_gets_only_its_fields(self):
"""Payload narrowing is shared with ``invoke_hook``: a callback declaring only ``event``
must not receive the additive ``gateway`` / ``telemetry_schema_version`` fields."""
def test_narrow_legacy_signature_still_gets_only_its_fields(self, caplog):
"""Payload narrowing and failure isolation are shared with ``invoke_hook``: a callback
declaring only ``event`` must not receive the additive ``gateway`` /
``telemetry_schema_version`` fields, and a raising callback is reported once and skipped
without losing its siblings' results."""
import asyncio
mgr = PluginManager()
@@ -2921,9 +2900,14 @@ class TestAsyncHookOnCallerLoop:
def narrow(event):
return {"seen": event}
async def boom(**_kw):
raise RuntimeError("async plugin blew up")
async def narrow_async(event):
return {"seen_async": event}
mgr._hooks.setdefault("pre_gateway_dispatch", []).extend([narrow, narrow_async])
results = asyncio.run(mgr.ainvoke_hook("pre_gateway_dispatch", event="e", gateway="g"))
mgr._hooks.setdefault("pre_gateway_dispatch", []).extend([narrow, boom, narrow_async])
with caplog.at_level(logging.WARNING, logger="hermes_cli.plugins"):
results = asyncio.run(mgr.ainvoke_hook("pre_gateway_dispatch", event="e", gateway="g"))
assert results == [{"seen": "e"}, {"seen_async": "e"}]
assert "async plugin blew up" in caplog.text

View File

@@ -1214,6 +1214,8 @@ def my_callback(event, gateway, session_store, **kwargs):
**Return value:** `None` or a dict. The first recognized action dict wins; remaining plugin results are ignored. Exceptions in plugin callbacks are caught and logged; the gateway always falls through to normal dispatch on error.
Callbacks may be `async def`: they are awaited on the gateway's own event loop, so awaiting loop-bound work (an `asyncio.Event`, an aiohttp session, `asyncio.to_thread`) makes progress and other inbound messages keep flowing while the callback runs. The hook is intentionally not bounded by `plugins.hook_callback_timeout` — dropping or passing a message on timeout are both wrong for a policy gate — so a callback that never returns holds up dispatch of that message.
| Return | Effect |
|--------|--------|
| `{"action": "skip", "reason": "..."}` | Drop the message — no agent reply, no pairing flow, no auth. Plugin is assumed to have handled it (e.g. silent-ingested into the transcript). |