diff --git a/hermes_cli/plugins.py b/hermes_cli/plugins.py index 954117cc8fdcb..c27f02c212e40 100644 --- a/hermes_cli/plugins.py +++ b/hermes_cli/plugins.py @@ -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 (``:`` or - # ``hermes:``) → 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 (``:`` or - ``hermes:``) 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 + (``:`` or ``hermes:``) 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]: diff --git a/tests/hermes_cli/test_plugin_event_bus.py b/tests/hermes_cli/test_plugin_event_bus.py index 6e4396716c3e0..6ea0e2aa578cb 100644 --- a/tests/hermes_cli/test_plugin_event_bus.py +++ b/tests/hermes_cli/test_plugin_event_bus.py @@ -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