fix(plugins): bound event delivery and own subscriptions

This commit is contained in:
Hans 2026-07-18 13:35:32 +08:00 committed by Teknium
parent 17030939db
commit 67168a391f
2 changed files with 434 additions and 52 deletions

View File

@ -34,6 +34,7 @@ so plugin-defined tools appear alongside the built-in tools.
from __future__ import annotations
import asyncio
import copy
import hashlib
import importlib.metadata
import importlib.util
@ -41,6 +42,7 @@ import inspect
import json
import logging
import os
import queue
import re
import sys
import threading
@ -291,6 +293,11 @@ HERMES_EVENT_NAMESPACE = "hermes"
# forever. When exceeded the over-deep emit is dropped (with a warning), not
# raised, so delivery always terminates cleanly.
_EVENT_EMIT_DEPTH_CAP = 8
# Maximum number of queued + currently-running events per manager generation.
# ``emit`` never waits for capacity: a full budget drops the new event with a
# warning so a blocked subscriber cannot back-pressure the emitter forever.
_EVENT_PENDING_CAP = 64
_EVENT_WORKER_STOP = object()
_NS_PARENT = "hermes_plugins"
@ -752,6 +759,25 @@ class RenderedPluginSystemPromptSection:
plugin: str
@dataclass(frozen=True)
class _EventSubscription:
"""Host-owned subscription ledger entry."""
owner: str
callback: Callable
@dataclass(frozen=True)
class _QueuedPluginEvent:
"""Immutable dispatch envelope consumed by the event worker."""
event: str
payload: Dict[str, Any]
subscriptions: tuple[_EventSubscription, ...]
depth: int
generation: int
@dataclass
class LoadedPlugin:
"""Runtime state for a single loaded plugin."""
@ -2174,13 +2200,15 @@ class PluginContext:
``ValueError`` and a logged warning fail-closed. The ``hermes:``
prefix is reserved for core.
Subscribers are invoked in registration order, each isolated in its
own ``try/except`` (one raising subscriber does not stop delivery to
the others). ``payload`` is passed to each callback as keyword
arguments (``cb(**payload)``), mirroring the hook convention.
Delivery is fire-and-forget through a host-owned, single-worker queue:
registration order is preserved, while a blocking subscriber cannot
stall the emitter. The queue has a bounded pending budget; a full
budget drops the new event with a warning. Each subscriber receives a
deep-copied payload and is isolated in its own ``try/except``. Awaitable
results are resolved through the existing loop-safe plugin path.
Returns the count of subscriber callbacks invoked (0 when there are
no subscribers, or when the recursion cap dropped the emit).
Returns the count of subscriber callbacks scheduled (0 when there are
no subscribers, or when the pending/recursion budget drops the emit).
"""
plugin_key = self.manifest.key or self.manifest.name
if not event or not isinstance(event, str):
@ -2204,6 +2232,10 @@ class PluginContext:
f"bare event name; the namespace is forced to '{plugin_key}:' "
f"and the '{HERMES_EVENT_NAMESPACE}:' prefix is reserved for core"
)
if payload is not None and not isinstance(payload, dict):
raise TypeError(
f"Plugin '{plugin_key}' emit() payload must be a dict or None"
)
full_event = f"{plugin_key}:{event}"
return self._manager._dispatch_event(full_event, payload or {})
@ -2214,15 +2246,17 @@ class PluginContext:
if core ever emits). Subscribing is unrestricted any plugin may
listen to any published event; only *emitting* is namespace-gated.
Callbacks are stored in registration order and invoked with the
emitter's payload as keyword arguments.
Callbacks are stored in registration order as host-owned ledger
entries. The owner key lets plugin unload/reload remove subscriptions
before any later event can invoke a zombie callback.
"""
if not event or not isinstance(event, str):
raise ValueError(
f"Plugin '{self.manifest.name}' subscribe() requires a "
f"non-empty event name"
)
self._manager._subscriptions.setdefault(event, []).append(callback)
plugin_key = self.manifest.key or self.manifest.name
self._manager._subscribe_event(plugin_key, event, callback)
logger.debug(
"Plugin %s subscribed to event: %s", self.manifest.name, event,
)
@ -2329,12 +2363,20 @@ class PluginManager:
self._aux_tasks: Dict[str, Dict[str, Any]] = {}
# Explicitly-selected, profile-scoped human approval transports.
self._approval_transports: Dict[str, Any] = {}
# Inter-plugin event bus: full event name (``<plugin_key>:<event>`` or
# ``hermes:<event>``) → list of subscriber callbacks in registration
# order. See PluginContext.subscribe/emit and _dispatch_event.
self._subscriptions: Dict[str, List[Callable]] = {}
# Per-thread re-entrancy depth for event dispatch, used to cap
# mutually-emitting plugins (see _EVENT_EMIT_DEPTH_CAP).
# Inter-plugin event bus. Subscriptions are owner-tagged ledger entries
# so unload/reload can remove zombie callbacks. A single daemon worker
# preserves registration order while keeping emitters non-blocking.
self._subscriptions: Dict[str, List[_EventSubscription]] = {}
self._event_lock = threading.RLock()
self._event_idle = threading.Condition(self._event_lock)
self._event_generation = 0
self._event_pending_by_generation: Dict[int, int] = {0: 0}
self._event_queue: queue.Queue[Any] = queue.Queue(
maxsize=_EVENT_PENDING_CAP
)
self._event_worker: Optional[threading.Thread] = None
# Per-worker chain depth caps mutually-emitting plugins even though each
# re-entrant emit is queued rather than invoked recursively.
self._emit_depth = threading.local()
# Slack Block Kit action handlers registered by plugins. Each entry
# is (matcher, callback, plugin_name); the Slack adapter wires them
@ -2408,7 +2450,7 @@ class PluginManager:
self._portable_mcp_servers.clear()
self._aux_tasks.clear()
self._approval_transports.clear()
self._subscriptions.clear()
self._reset_event_bus()
self._slack_action_handlers.clear()
self._context_engine = None
# Set the flag up front as a re-entrancy guard (a plugin's register()
@ -3219,6 +3261,10 @@ class PluginManager:
except Exception as exc:
loaded.error = str(exc)
# register() may have subscribed before raising. Remove those
# owner-tagged entries so a failed/unloaded plugin cannot leave a
# callable reachable from later event dispatch.
self._remove_plugin_subscriptions(_plugin_id)
logger.warning(
"Failed to load plugin '%s': %s",
manifest.name, exc, exc_info=_PLUGINS_DEBUG,
@ -3410,18 +3456,158 @@ class PluginManager:
)
return results
def _subscribe_event(
self,
owner: str,
event: str,
callback: Callable,
) -> None:
"""Add an owner-tagged event subscription in registration order."""
if not callable(callback):
raise TypeError("Event subscriber callback must be callable")
entry = _EventSubscription(owner=owner, callback=callback)
with self._event_lock:
self._subscriptions.setdefault(event, []).append(entry)
def _remove_plugin_subscriptions(self, owner: str) -> int:
"""Remove every subscription owned by *owner* and return the count.
Queued dispatch envelopes re-check ledger membership before each
callback, so removing an owner also cancels callbacks already snapshotted
by an event that has not reached that subscriber yet.
TODO(#64229): when the central plugin ownership ledger / registration
handles land, route this owner-tagged bookkeeping through that ledger
so per-plugin unload cancels event subscriptions alongside every other
registration surface. This method is the integration seam.
"""
removed = 0
with self._event_lock:
for event in list(self._subscriptions):
entries = self._subscriptions[event]
retained = [entry for entry in entries if entry.owner != owner]
removed += len(entries) - len(retained)
if retained:
self._subscriptions[event] = retained
else:
del self._subscriptions[event]
return removed
def _reset_event_bus(self) -> None:
"""Cancel the current event generation and clear its subscriptions."""
with self._event_lock:
old_queue = self._event_queue
had_worker = self._event_worker is not None
self._event_generation += 1
self._subscriptions.clear()
self._event_queue = queue.Queue(maxsize=_EVENT_PENDING_CAP)
self._event_worker = None
self._event_pending_by_generation.setdefault(
self._event_generation, 0
)
# Drop work that has not started. A currently-running callback
# cannot be force-killed safely, but generation + ledger checks stop
# it before the next subscriber and prevent all queued callbacks.
while True:
try:
item = old_queue.get_nowait()
except queue.Empty:
break
try:
if item is not _EVENT_WORKER_STOP:
self._mark_event_done(item.generation)
finally:
old_queue.task_done()
if had_worker:
old_queue.put_nowait(_EVENT_WORKER_STOP)
self._event_idle.notify_all()
def _ensure_event_worker_locked(self) -> None:
worker = self._event_worker
if worker is not None and worker.is_alive():
return
dispatch_queue = self._event_queue
worker = threading.Thread(
target=self._event_worker_loop,
args=(dispatch_queue,),
name="hermes-plugin-events",
daemon=True,
)
self._event_worker = worker
worker.start()
def _event_worker_loop(self, dispatch_queue: queue.Queue[Any]) -> None:
while True:
item = dispatch_queue.get()
try:
if item is _EVENT_WORKER_STOP:
return
self._deliver_event(item)
finally:
if item is not _EVENT_WORKER_STOP:
self._mark_event_done(item.generation)
dispatch_queue.task_done()
def _mark_event_done(self, generation: int) -> None:
with self._event_idle:
pending = self._event_pending_by_generation.get(generation, 0)
if pending > 0:
self._event_pending_by_generation[generation] = pending - 1
self._event_idle.notify_all()
def _deliver_event(self, item: _QueuedPluginEvent) -> None:
"""Deliver one queued event on the host-owned worker thread."""
with self._event_lock:
if item.generation != self._event_generation:
return
previous_depth = getattr(self._emit_depth, "value", 0)
self._emit_depth.value = item.depth
try:
for subscription in item.subscriptions:
with self._event_lock:
if item.generation != self._event_generation:
break
# Owner unload may remove this exact ledger entry after the
# event was queued but before its callback starts.
if not any(
current is subscription
for current in self._subscriptions.get(item.event, [])
):
continue
callback = subscription.callback
try:
# A fresh deep copy per subscriber prevents one callback
# from mutating the emitter's nested values or the payload
# observed by the next subscriber.
owned_payload = copy.deepcopy(item.payload)
result = callback(**owned_payload)
resolve_plugin_command_result(result)
except Exception as exc:
logger.warning(
"Event '%s' subscriber %s raised: %s",
item.event,
getattr(callback, "__name__", repr(callback)),
exc,
)
finally:
self._emit_depth.value = previous_depth
def _wait_for_event_dispatch(self, timeout: float = 2.0) -> bool:
"""Wait for the current event generation to become idle (test helper)."""
with self._event_idle:
generation = self._event_generation
return self._event_idle.wait_for(
lambda: self._event_pending_by_generation.get(generation, 0) == 0,
timeout=timeout,
)
def _dispatch_event(self, event: str, payload: Dict[str, Any]) -> int:
"""Deliver *event* to its subscribers; return the number invoked.
"""Queue *event* without blocking; return subscriber count scheduled.
Mirrors :meth:`invoke_hook`: iterate subscribers in registration
order, isolate each in its own ``try/except`` so one raising
subscriber cannot break delivery to the rest, and pass the payload
as keyword arguments.
A subscriber may itself call ``ctx.emit`` (re-entrant dispatch). A
per-thread depth counter caps recursion at
:data:`_EVENT_EMIT_DEPTH_CAP`; over-deep emits are dropped with a
single warning (never raised) so delivery always terminates.
A single daemon worker preserves registration order. Pending work is
bounded per manager generation so a blocking subscriber can consume at
most one worker while later emits are dropped once the budget is full.
"""
depth = getattr(self._emit_depth, "value", 0)
if depth >= _EVENT_EMIT_DEPTH_CAP:
@ -3431,26 +3617,41 @@ class PluginManager:
_EVENT_EMIT_DEPTH_CAP, event,
)
return 0
callbacks = list(self._subscriptions.get(event, []))
if not callbacks:
return 0
self._emit_depth.value = depth + 1
invoked = 0
try:
for cb in callbacks:
invoked += 1
try:
cb(**payload)
except Exception as exc:
logger.warning(
"Event '%s' subscriber %s raised: %s",
event,
getattr(cb, "__name__", repr(cb)),
exc,
)
finally:
self._emit_depth.value = depth
return invoked
with self._event_lock:
subscriptions = tuple(self._subscriptions.get(event, []))
if not subscriptions:
return 0
generation = self._event_generation
pending = self._event_pending_by_generation.get(generation, 0)
if pending >= _EVENT_PENDING_CAP:
logger.warning(
"Event bus pending budget (%d) exhausted while dispatching "
"'%s' — dropping this emit",
_EVENT_PENDING_CAP,
event,
)
return 0
item = _QueuedPluginEvent(
event=event,
payload=dict(payload),
subscriptions=subscriptions,
depth=depth + 1,
generation=generation,
)
try:
self._event_queue.put_nowait(item)
except queue.Full:
logger.warning(
"Event bus pending budget (%d) exhausted while dispatching "
"'%s' — dropping this emit",
_EVENT_PENDING_CAP,
event,
)
return 0
self._event_pending_by_generation[generation] = pending + 1
self._ensure_event_worker_locked()
return len(subscriptions)
def has_hook(self, hook_name: str) -> bool:
"""Return True when at least one callback is registered for a hook."""
@ -4245,13 +4446,17 @@ def get_plugin_auxiliary_tasks() -> List[Dict[str, Any]]:
def get_plugin_subscriptions() -> Dict[str, List[Callable]]:
"""Return the inter-plugin event bus subscription registry.
Maps each fully-qualified event name (``<plugin_key>:<event>`` or
``hermes:<event>``) to its list of subscriber callbacks in registration
order. Triggers idempotent plugin discovery so callers can read the
registry before any explicit ``discover_plugins()`` call.
Returns a snapshot mapping each fully-qualified event name
(``<plugin_key>:<event>`` or ``hermes:<event>``) to subscriber callbacks in
registration order. Owner ledger metadata stays private to the manager.
Triggers idempotent plugin discovery before reading the snapshot.
"""
manager = _ensure_plugins_discovered()
return manager._subscriptions
with manager._event_lock:
return {
event: [entry.callback for entry in entries]
for event, entries in manager._subscriptions.items()
}
def get_plugin_toolsets() -> List[tuple]:

View File

@ -4,7 +4,10 @@ Covers:
- Two plugins communicate via emit/subscribe; emit returns listener count
- Namespace is FORCED to the emitting plugin's own key
- Namespace spoofing (hermes:, foreign, already-colon'd) is rejected
- Per-callback isolation: one raising subscriber does not break the rest
- Non-blocking bounded delivery for synchronous subscribers
- Per-callback isolation and deep-copied payload ownership
- Async subscribers resolved through the loop-safe host path
- Owner unload / generation reset cancel zombie callbacks
- Recursion cap: mutually-emitting plugins terminate + warn
- Manifest emits/listens parsed as optional advisory fields
- `hermes plugins show` output includes emits/listens
@ -12,7 +15,9 @@ Covers:
from __future__ import annotations
import asyncio
import logging
import threading
import pytest
@ -40,6 +45,10 @@ def _fresh_manager() -> PluginManager:
return manager
def _drain(manager: PluginManager) -> None:
assert manager._wait_for_event_dispatch(timeout=2.0)
# ── 1. Two plugins communicate ───────────────────────────────────────────────
@ -56,6 +65,7 @@ def test_two_plugins_communicate():
# A subscribes to b:ping; B emits the bare name "ping".
ctx_a.subscribe("b:ping", on_ping)
count = ctx_b.emit("ping", {"n": 42})
_drain(manager)
assert count == 1 # one listener invoked
assert received == [{"n": 42}]
@ -75,6 +85,7 @@ def test_emit_none_payload_delivers_empty_kwargs():
seen = []
ctx_a.subscribe("b:ping", lambda **p: seen.append(p))
count = ctx_b.emit("ping") # payload omitted
_drain(manager)
assert count == 1
assert seen == [{}]
@ -96,6 +107,7 @@ def test_namespace_forced_to_emitter_key():
ctx_a.subscribe("a:ping", lambda **p: delivered_events.append("a:ping"))
ctx_b.emit("ping")
_drain(manager)
# Delivered under the emitter's own key ("b"), never "a".
assert delivered_events == ["b:ping"]
@ -110,6 +122,7 @@ def test_namespace_falls_back_to_name_when_key_empty():
got = []
ctx.subscribe("plugin_named:evt", lambda **p: got.append(p))
count = ctx.emit("evt", {"v": 1})
_drain(manager)
assert count == 1
assert got == [{"v": 1}]
@ -184,6 +197,7 @@ def test_per_callback_isolation(caplog):
with caplog.at_level(logging.WARNING):
count = ctx_b.emit("ping", {"ok": True})
_drain(manager)
# Both listeners were invoked despite the first raising.
assert count == 2
@ -192,6 +206,168 @@ def test_per_callback_isolation(caplog):
for r in caplog.records)
def test_emit_returns_before_blocking_subscriber_finishes():
manager = _fresh_manager()
ctx_a = _make_ctx(manager, "plugin_a", key="a")
ctx_b = _make_ctx(manager, "plugin_b", key="b")
entered = threading.Event()
release = threading.Event()
emit_returned = threading.Event()
result = {}
def blocking(**payload):
entered.set()
release.wait(timeout=2.0)
def call_emit():
result["count"] = ctx_b.emit("ping")
emit_returned.set()
ctx_a.subscribe("b:ping", blocking)
emitter = threading.Thread(target=call_emit)
emitter.start()
try:
assert emit_returned.wait(timeout=1.0)
assert entered.wait(timeout=1.0)
finally:
release.set()
emitter.join(timeout=2.0)
_drain(manager)
assert result["count"] == 1
def test_pending_budget_drops_new_event_without_blocking(monkeypatch, caplog):
from hermes_cli import plugins as plugins_mod
monkeypatch.setattr(plugins_mod, "_EVENT_PENDING_CAP", 1)
manager = _fresh_manager()
ctx_a = _make_ctx(manager, "plugin_a", key="a")
ctx_b = _make_ctx(manager, "plugin_b", key="b")
entered = threading.Event()
release = threading.Event()
def blocking(**payload):
entered.set()
release.wait(timeout=2.0)
ctx_a.subscribe("b:ping", blocking)
assert ctx_b.emit("ping") == 1
assert entered.wait(timeout=1.0)
try:
with caplog.at_level(logging.WARNING):
assert ctx_b.emit("ping") == 0
finally:
release.set()
_drain(manager)
assert "pending budget" in caplog.text
def test_each_subscriber_receives_deep_copied_payload():
manager = _fresh_manager()
ctx_a = _make_ctx(manager, "plugin_a", key="a")
ctx_b = _make_ctx(manager, "plugin_b", key="b")
original = {"nested": {"value": 1}}
observed = []
def mutate(**payload):
payload["nested"]["value"] = 99
def observe(**payload):
observed.append(payload["nested"]["value"])
ctx_a.subscribe("b:ping", mutate)
ctx_a.subscribe("b:ping", observe)
assert ctx_b.emit("ping", original) == 2
_drain(manager)
assert original == {"nested": {"value": 1}}
assert observed == [1]
def test_async_subscriber_is_awaited():
manager = _fresh_manager()
ctx_a = _make_ctx(manager, "plugin_a", key="a")
ctx_b = _make_ctx(manager, "plugin_b", key="b")
observed = []
async def on_ping(**payload):
await asyncio.sleep(0)
observed.append(payload["value"])
ctx_a.subscribe("b:ping", on_ping)
assert ctx_b.emit("ping", {"value": 7}) == 1
_drain(manager)
assert observed == [7]
def test_remove_plugin_subscriptions_cancels_owner_entries():
manager = _fresh_manager()
ctx_a = _make_ctx(manager, "plugin_a", key="a")
ctx_b = _make_ctx(manager, "plugin_b", key="b")
observed = []
ctx_a.subscribe("b:ping", lambda **payload: observed.append(payload))
manager._remove_plugin_subscriptions("a")
assert ctx_b.emit("ping", {"value": 1}) == 0
_drain(manager)
assert observed == []
assert "b:ping" not in manager._subscriptions
def test_owner_removal_cancels_callback_already_snapshotted_in_queue():
manager = _fresh_manager()
ctx_gate = _make_ctx(manager, "gate", key="gate")
ctx_a = _make_ctx(manager, "plugin_a", key="a")
ctx_b = _make_ctx(manager, "plugin_b", key="b")
entered = threading.Event()
release = threading.Event()
observed = []
def blocking(**payload):
entered.set()
release.wait(timeout=2.0)
ctx_gate.subscribe("b:ping", blocking)
ctx_a.subscribe("b:ping", lambda **payload: observed.append(payload))
assert ctx_b.emit("ping", {"value": 1}) == 2
assert entered.wait(timeout=1.0)
manager._remove_plugin_subscriptions("a")
release.set()
_drain(manager)
assert observed == []
def test_event_bus_reset_cancels_queued_generation():
manager = _fresh_manager()
ctx_gate = _make_ctx(manager, "gate", key="gate")
ctx_a = _make_ctx(manager, "plugin_a", key="a")
ctx_b = _make_ctx(manager, "plugin_b", key="b")
entered = threading.Event()
release = threading.Event()
observed = []
def blocking(**payload):
entered.set()
release.wait(timeout=2.0)
ctx_gate.subscribe("b:ping", blocking)
ctx_a.subscribe("b:ping", lambda **payload: observed.append(payload))
assert ctx_b.emit("ping", {"value": 1}) == 2
assert entered.wait(timeout=1.0)
old_worker = manager._event_worker
manager._reset_event_bus()
release.set()
assert old_worker is not None
old_worker.join(timeout=2.0)
assert not old_worker.is_alive()
assert observed == []
assert manager._subscriptions == {}
# ── 5. Recursion cap ─────────────────────────────────────────────────────────
@ -217,6 +393,7 @@ def test_recursion_cap_terminates(caplog):
with caplog.at_level(logging.WARNING):
# Kick off the loop — must terminate, not hang or RecursionError.
result = ctx_b.emit("ping")
_drain(manager)
# Returned cleanly.
assert result == 1