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:
@@ -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, ...]:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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). |
|
||||
|
||||
Reference in New Issue
Block a user