diff --git a/hermes_cli/plugins.py b/hermes_cli/plugins.py index 6eb97cf7c1..09f83f8c53 100644 --- a/hermes_cli/plugins.py +++ b/hermes_cli/plugins.py @@ -11,7 +11,6 @@ and an ``__init__.py`` exposing ``register(ctx)``. Plugins register callbacks fo from __future__ import annotations import asyncio -import builtins import importlib.metadata import inspect import json @@ -24,84 +23,44 @@ import threading import types from contextlib import suppress from dataclasses import dataclass, field +from functools import cached_property from pathlib import Path from typing import Any, Callable, Dict, List, Mapping, Optional, Set, Tuple, Union from hermes_constants import get_hermes_home, hermes_home_key from registration_lifecycle import replacement_coordinator from utils import env_var_enabled -from hermes_cli.config import cfg_get, load_config_readonly +from hermes_cli.config import load_config_readonly from hermes_cli.middleware import VALID_MIDDLEWARE -from hermes_cli.plugin_capabilities import ( # noqa: F401 — re-exported - CAPABILITY_REGISTRY, - VALID_CAPABILITY_IDS, - plugin_capability_granted, -) +from hermes_cli.plugin_capabilities import plugin_capability_granted from hermes_cli.relay_plugin_cutover import RELAY_PLUGINS_CONFIG_ENV, legacy_relay_plugin_keys +# Sibling modules' names are re-exported here (origin) so plugins and tests keep one import path. from hermes_cli.plugins_manifest import ( # noqa: F401 — re-exported - manifest_key, - parse_manifest_file, - portable_plugin_manifest, - _CONFIG_SCHEMA_TYPES, - SUPPORTED_MANIFEST_VERSION, - PluginManifest, - _portable_skill_namespace, - resolve_module_origin, - resolve_plugin_load_order, + _CONFIG_SCHEMA_TYPES, SUPPORTED_MANIFEST_VERSION, PluginManifest, _portable_skill_namespace, + manifest_key, parse_manifest_file, resolve_module_origin, resolve_plugin_load_order, validate_config_schema, ) from hermes_cli.plugins_discovery import ( # noqa: F401 — re-exported - collect_directory_manifests, - gate_manifest, - scan_directory, - ENTRY_POINT_CAPABILITIES_GROUP, - ENTRY_POINTS_GROUP, - _get_disabled_plugins, - _get_enabled_plugins, - discover_entrypoint_manifests, + ENTRY_POINTS_GROUP, _get_disabled_plugins, _get_enabled_plugins, collect_directory_manifests, + discover_entrypoint_manifests, gate_manifest, scan_directory, ) -from hermes_cli.plugins_loader import ( # noqa: F401 — re-exported - PluginLoaderMixin, - _NS_PARENT, - _MODULE_NAMESPACE_LOCK, - _BARE_MODULE_SCOPE, - _evict_modules, - _serialized_replacement, - _plugin_home_scope, +from hermes_cli.plugins_loader import ( + PluginLoaderMixin, _BARE_MODULE_SCOPE, _MODULE_NAMESPACE_LOCK, _NS_PARENT, _evict_modules, + _plugin_home_scope, _serialized_replacement, ) from hermes_cli.plugins_dispatch import ( # noqa: F401 — re-exported - PluginDispatchMixin, - _HOOK_TIMEOUT_SUPPRESSION_SECONDS, - _PRE_TOOL_CALL_TIMEOUT_BLOCK_MESSAGE, - SYSTEM_PROMPT_SECTION_POSITIONS, - DEFAULT_SYSTEM_PROMPT_SECTION_MAX_CHARS, - MAX_SYSTEM_PROMPT_SECTION_CHARS, - MAX_SYSTEM_PROMPT_SECTIONS, - MAX_SYSTEM_PROMPT_SECTIONS_TOTAL_CHARS, - PLUGIN_SECTIONS_START, - PLUGIN_SECTIONS_END, + DEFAULT_SYSTEM_PROMPT_SECTION_MAX_CHARS, HERMES_EVENT_NAMESPACE, MAX_SYSTEM_PROMPT_SECTION_CHARS, + MAX_SYSTEM_PROMPT_SECTIONS_TOTAL_CHARS, PLUGIN_SECTIONS_END, PLUGIN_SECTIONS_START, + SYSTEM_PROMPT_SECTION_POSITIONS, _EVENT_EMIT_DEPTH_CAP, _EVENT_PENDING_CAP, + _HOOK_CALLBACK_TIMEOUT_SECS, _HOOK_TIMEOUT_SUPPRESSION_SECONDS, _MAX_HOOK_CALLBACK_TIMEOUT_SECS, + _PRE_TOOL_CALL_TIMEOUT_BLOCK_MESSAGE, PluginDispatchMixin, PluginSystemPromptSection, + RenderedPluginSystemPromptSection, _EventSubscription, format_system_prompt_sections, is_valid_system_prompt_section_id, - format_system_prompt_section, - format_system_prompt_sections, - HERMES_EVENT_NAMESPACE, - _EVENT_EMIT_DEPTH_CAP, - _EVENT_PENDING_CAP, - PluginSystemPromptSection, - RenderedPluginSystemPromptSection, - _EventSubscription, - _HOOK_CALLBACK_TIMEOUT_SECS, - _MAX_HOOK_CALLBACK_TIMEOUT_SECS, ) -from hermes_cli.plugins_ledger import ( # noqa: F401 — re-exported - PluginLedgerMixin, - PluginRegistration, -) -from hermes_cli.plugins_state import ( # noqa: F401 — re-exported - PluginState, - _locked_plugin_state, - _nested_plugin_mapping, - _nested_plugin_value, - _plugin_relative_segments, +from hermes_cli.plugins_ledger import PluginLedgerMixin, PluginRegistration +from hermes_cli.plugins_state import ( + PluginState, _locked_plugin_state, _nested_plugin_mapping, _nested_plugin_value, + _plugin_relative_segments, _plugin_settings_entry, ) @@ -120,10 +79,9 @@ class PluginToolOverrideError(PermissionError): logger = logging.getLogger(__name__) - # ``HERMES_PLUGINS_DEBUG=1`` tees verbose discovery logs to stderr in addition to agent.log. Read # once at import; tests flip it mid-process via ``_install_plugin_debug_handler(force=True)``. -_PLUGINS_DEBUG = os.getenv("HERMES_PLUGINS_DEBUG", "").strip().lower() in {"1", "true", "yes", "on"} +_PLUGINS_DEBUG = env_var_enabled("HERMES_PLUGINS_DEBUG") _DEBUG_HANDLER_INSTALLED = False @@ -131,7 +89,7 @@ def _install_plugin_debug_handler(force: bool = False) -> None: """When HERMES_PLUGINS_DEBUG is on, tee plugin logs to stderr at DEBUG (once per process).""" global _DEBUG_HANDLER_INSTALLED, _PLUGINS_DEBUG if force: - _PLUGINS_DEBUG = os.getenv("HERMES_PLUGINS_DEBUG", "").strip().lower() in {"1", "true", "yes", "on"} + _PLUGINS_DEBUG = env_var_enabled("HERMES_PLUGINS_DEBUG") if not _PLUGINS_DEBUG or _DEBUG_HANDLER_INSTALLED: return handler = logging.StreamHandler(sys.stderr) @@ -148,71 +106,67 @@ _install_plugin_debug_handler() VALID_HOOKS: Set[str] = { "pre_tool_call", "post_tool_call", "transform_terminal_output", "transform_tool_result", - # Return a string to replace the response text (first non-None wins) or None to leave it. + # transform_llm_output: return a replacement string (first non-None wins) or None. "transform_llm_output", "pre_llm_call", "post_llm_call", - # Streaming observers fired off the token path by agent.plugin_stream_hooks; payloads are - # immutable normalized text/lifecycle and cannot transform the stream. + # Streaming observers (agent.plugin_stream_hooks), off the token path; payloads are immutable + # normalized text/lifecycle and cannot transform the stream. "on_stream_start", "on_stream_delta", "on_stream_end", "on_interim_message", - # Fired once per turn when the agent edited code and is about to verify/finish. Return - # {"action": "continue", "message"} (or Claude-Code Stop shape {"decision": "block", "reason"}) - # to keep going; anything else finishes. Bounded by agent.max_verify_nudges. + # pre_verify: once per turn when the agent edited code and is about to verify/finish. Return + # {"action": "continue", "message"} (or Claude-Code Stop {"decision": "block", "reason"}) to keep + # going; anything else finishes. Bounded by agent.max_verify_nudges. "pre_verify", "pre_api_request", "post_api_request", "api_request_error", - # Fired once per failed API call BEFORE agent/error_classifier.classify_api_error(). Kwargs: - # provider, model, status_code, error_type, error_code, error_message, error_body, error, - # approx_tokens, context_length, num_messages. Return None, or {"reason": - # (required), "retryable"/"should_compress"/"should_rotate_credential"/"should_fallback": bool, - # "message": str, "error_context": dict}. Run-all-then-pick-first (see - # get_plugin_error_classification). Privacy: error_message/error_body may be unredacted. + # transform_api_error_classification: once per failed API call BEFORE + # agent/error_classifier.classify_api_error(). Kwargs: provider, model, status_code, error_type, + # error_code, error_message, error_body, error, approx_tokens, context_length, num_messages. + # Return None or {"reason": (required), "retryable"/"should_compress"/ + # "should_rotate_credential"/"should_fallback": bool, "message": str, "error_context": dict}. + # Run-all-then-pick-first (see get_plugin_error_classification). Privacy: error_message/ + # error_body may be unredacted. "transform_api_error_classification", "on_session_start", "on_session_end", "on_session_finalize", "on_session_reset", - # Successful skill lifecycle facts (local skill name visible to plugins). + # on_skill_lifecycle: successful skill lifecycle facts (local skill name visible to plugins). "on_skill_lifecycle", "subagent_start", "subagent_stop", - # Once per incoming MessageEvent, after the internal-event guard, BEFORE auth/pairing and - # dispatch. Kwargs: event, gateway, session_store. Return {"action": "skip", "reason"} -> drop; - # {"action": "rewrite", "text"} -> replace event.text; {"action": "allow"} / None -> normal. + # pre_gateway_dispatch: once per incoming MessageEvent, after the internal-event guard, BEFORE + # auth/pairing and dispatch. Kwargs: event, gateway, session_store. Return {"action": "skip", + # "reason"} -> drop; {"action": "rewrite", "text"} -> replace event.text; "allow"/None -> normal. "pre_gateway_dispatch", - # Approval observers (tools/approval.py). Return values ignored — plugins cannot veto or - # pre-answer (use pre_tool_call). Kwargs: command, description, pattern_key, pattern_keys, - # session_key, surface: "cli" | "gateway" | "smart"; post_approval_response adds choice - # ("once"|"session"|"always"|"deny"|"timeout"|"smart_approve"|"smart_deny") and decided_by. + # Approval observers (tools/approval.py); returns ignored — plugins cannot veto or pre-answer + # (use pre_tool_call). Kwargs: command, description, pattern_key, pattern_keys, session_key, + # surface: "cli"|"gateway"|"smart"; post_approval_response adds choice ("once"|"session"| + # "always"|"deny"|"timeout"|"smart_approve"|"smart_deny") and decided_by. "pre_approval_request", "post_approval_response", - # Fired by transcribe_audio after provider resolution, BEFORE any backend runs. Kwargs: - # file_path, provider, model, language, prompt, source. Return None or a dict mutating - # prompt/language/model (registration order, last-writer-wins; file_path is read-only). + # pre_transcription: after provider resolution, BEFORE any backend runs. Kwargs: file_path, + # provider, model, language, prompt, source. Return None or a dict mutating prompt/language/ + # model (registration order, last-writer-wins; file_path is read-only). "pre_transcription", # Kanban task observers (hermes_cli.kanban_db), fired AFTER the DB commit so a slow plugin never - # holds the SQLite write lock. Return values ignored. Process matters: claimed fires in the - # DISPATCHER right before spawn; completed/blocked fire in the WORKER (or whichever process - # drove it). Kwargs: task_id, board, assignee, run_id, profile_name; completed adds summary, - # blocked adds reason. + # holds the SQLite write lock; returns ignored. claimed fires in the DISPATCHER right before + # spawn; completed/blocked fire in the WORKER (or whichever process drove it). Kwargs: task_id, + # board, assignee, run_id, profile_name; completed adds summary, blocked adds reason. "kanban_task_claimed", "kanban_task_completed", "kanban_task_blocked", - # Kanban worker/mutation/tick observers; return values ignored; fire sites short-circuit on - # has_hook(). Kwargs: task_id, profile_name, board, assignee, run_id plus: + # Kanban worker/mutation/tick observers; returns ignored; fire sites short-circuit on + # has_hook(). Kwargs: task_id, profile_name, board, assignee, run_id plus, per hook: # worker_spawned (DISPATCHER, after PID persisted, inside the dispatch lock — stay fast): # worker_pid, workspace_path (privacy: project layout/usernames). - "on_kanban_worker_spawned", # worker_exited (tick-derived on dead-PID reclaim): worker_pid, exit_kind ("clean_exit" | # "rate_limited" | "nonzero_exit" | "signaled" | "unknown"), exit_code, outcome, retry_status. - "on_kanban_worker_exited", # worker_stale_claim (TTL-expired claim reclaimed; live-PID extensions do NOT fire): # worker_pid, heartbeat_stale, retry_status. - "on_kanban_worker_stale_claim", # task_updated (committed task-row write outside claim/complete/block, in whichever process # committed it): changed_fields — field NAMES only, never values. - "on_kanban_task_updated", - # dispatch_tick: once per dispatch_once, strictly AFTER the dispatch lock is released. Kwargs: - # board, profile_name, dry_run, outcome ("ok"|"skipped_locked"|"idle"), result: DispatchResult + # dispatch_tick (once per dispatch_once, strictly AFTER the dispatch lock is released): board, + # profile_name, dry_run, outcome ("ok"|"skipped_locked"|"idle"), result: DispatchResult # (privacy: task ids, assignees, workspace paths). - "on_kanban_dispatch_tick", - # Gateway platform-boundary observer: normalized envelopes only, never raw SDK objects or - # adapter handles. Kwargs: platform, event_type, payload (event_type-local; see hooks.md). - # New event types land only together with real fire-sites. + "on_kanban_worker_spawned", "on_kanban_worker_exited", "on_kanban_worker_stale_claim", + "on_kanban_task_updated", "on_kanban_dispatch_tick", + # gateway_platform_event: normalized envelopes only, never raw SDK objects or adapter handles. + # Kwargs: platform, event_type, payload (event_type-local; see hooks.md). New event types land + # only together with real fire-sites. "gateway_platform_event", - # Fired BEFORE a recognized slash command's handler on CLI and gateway canonical dispatch. - # Return values IGNORED in v1. Deliberately NOT fired for the gateway's running-agent intercept - # path (/stop, /approve, busy_policy) — a slow/hostile plugin must not touch the operator's - # escape hatches. Kwargs: surface, command (canonical), alias_used, args_raw, session_key, - # platform. + # pre_command: BEFORE a recognized slash command's handler on CLI and gateway canonical dispatch; + # returns IGNORED in v1. Deliberately NOT fired for the gateway's running-agent intercept path + # (/stop, /approve, busy_policy) — a slow/hostile plugin must not touch the operator's escape + # hatches. Kwargs: surface, command (canonical), alias_used, args_raw, session_key, platform. "pre_command", } @@ -221,6 +175,7 @@ VALID_HOOKS: Set[str] = { SHELL_UNSUPPORTED_HOOKS: Set[str] = {"transform_api_error_classification"} _env_enabled = env_var_enabled # imported by plugins/memory +_UNSET = object() @dataclass @@ -245,11 +200,7 @@ class PluginContext: def __init__(self, manifest: PluginManifest, manager: "PluginManager"): self.manifest = manifest self._manager = manager - # Lazy-built facades (see the matching properties). - self._llm: Any = None - self._subagent_lifecycle: Any = None - self._state: PluginState | None = None - self._platform_actions: Any = None + self._llm: Any = None # lazy; tests preseed it (see ``llm``) @property def plugin_id(self) -> str: @@ -264,20 +215,21 @@ class PluginContext: for key, loaded in self._manager._plugins.items() ) - def get_config(self, key: str, default: Any = None) -> Any: - """Read plugin-relative ``plugins.entries..settings.`` (falls back to the - legacy ``config`` subtree for migration compatibility).""" + def _segments(self, key: str) -> tuple[str, ...]: + """Validated plugin-relative settings path (warn + re-raise on rejection).""" try: - segments = _plugin_relative_segments(key) + return _plugin_relative_segments(key) except ValueError: logger.warning("Rejected config path %r from plugin %s", key, self.plugin_id) raise + + def get_config(self, key: str, default: Any = None) -> Any: + """Read plugin-relative ``plugins.entries..settings.`` (falls back to the + legacy ``config`` subtree for migration compatibility).""" + segments = self._segments(key) from hermes_cli.config import load_config_readonly - config = load_config_readonly() or {} - plugins = config.get("plugins") if isinstance(config, Mapping) else None - entries = plugins.get("entries") if isinstance(plugins, Mapping) else None - entry = entries.get(self.plugin_id) if isinstance(entries, Mapping) else None - if not isinstance(entry, Mapping): + entry = _plugin_settings_entry(load_config_readonly() or {}, self.plugin_id) + if entry is None: return default missing = object() value = _nested_plugin_value(entry.get("settings"), segments, missing) @@ -287,47 +239,37 @@ class PluginContext: def set_config(self, key: str, value: Any) -> None: """Atomically write one value in this plugin's ``settings`` subtree.""" - try: - segments = _plugin_relative_segments(key) - except ValueError: - logger.warning("Rejected config path %r from plugin %s", key, self.plugin_id) - raise + segments = self._segments(key) from hermes_cli import config as config_mod if config_mod.is_managed(): raise PermissionError("Plugin settings cannot be changed in a managed install") from hermes_cli import managed_scope - dotted_path = ".".join(("plugins", "entries", self.plugin_id, "settings", *segments)) + full_path = ("plugins", "entries", self.plugin_id, "settings", *segments) + dotted_path = ".".join(full_path) if managed_scope.is_key_managed(dotted_path): raise PermissionError(f"Plugin setting {dotted_path!r} is administrator-managed") - full_path = ("plugins", "entries", self.plugin_id, "settings", *segments) partial = _nested_plugin_mapping(full_path[:4], _nested_plugin_mapping(segments, value)) # The lock covers merge-read plus atomic save so sibling plugin writes (threads or # processes) cannot race between the two steps. - with _locked_plugin_state(config_mod.get_config_path()): - with config_mod._CONFIG_LOCK: - # Fail closed on malformed YAML: save_config degrades parse failures to {} — safe - # for reads, destructive for read-modify-write. - config_mod.read_user_config_raw() - config_mod.save_config(partial, preserve_keys={full_path}, merge_existing=True) + with _locked_plugin_state(config_mod.get_config_path()), config_mod._CONFIG_LOCK: + # Fail closed on malformed YAML: save_config degrades parse failures to {} — safe + # for reads, destructive for read-modify-write. + config_mod.read_user_config_raw() + config_mod.save_config(partial, preserve_keys={full_path}, merge_existing=True) - @property + @cached_property def state(self) -> PluginState: - """Return this plugin's profile-scoped durable JSON state facade.""" - if self._state is None: - self._state = PluginState(self.plugin_id, self.manifest.skill_namespace) - return self._state + """This plugin's profile-scoped durable JSON state facade.""" + return PluginState(self.plugin_id, self.manifest.skill_namespace) - @property + @cached_property def platform_actions(self): """Capability-gated platform action facade (``add_reaction``, ``set_thread_title``). Every call re-checks ``gateway.platform_actions`` (legacy ``plugins.entries..allow_platform_actions``, default OFF) and returns ``{"ok": bool, ...}`` — verbs never raise into hook dispatch; no adapter handles or raw SDK objects.""" - if self._platform_actions is None: - from hermes_cli.platform_actions import PlatformActions - - self._platform_actions = PlatformActions(self.plugin_id) - return self._platform_actions + from hermes_cli.platform_actions import PlatformActions + return PlatformActions(self.plugin_id) def _wrong_type(self, obj: Any, base_class: type, label: str, article: str = "a") -> bool: """Warn-and-ignore gate shared by every registrar that requires a base class.""" @@ -357,23 +299,21 @@ class PluginContext: restore: Callable[[Any], bool], ) -> PluginRegistration: """Track one generation in a replaceable manager-local registration slot.""" - lease = replacement_coordinator.acquire( - slot, current=current, previous=previous, restore=restore - ) + lease = replacement_coordinator.acquire(slot, current=current, previous=previous, restore=restore) return self._track(kind, key, lease.dispose) def _track_mapping_entry( - self, kind: str, key: str, mapping: Dict[str, Any], entry: Any, previous: Any + self, kind: str, key: str, mapping: Dict[str, Any], entry: Any, previous: Any = _UNSET, ) -> PluginRegistration: """Store ``entry`` under ``key`` in a manager-local mapping and lease the slot; unload restores - ``previous`` (or removes the key) only while ``entry`` is still current.""" + ``previous`` (default: the displaced entry, or removes the key) only while ``entry`` is still + current.""" + if previous is _UNSET: + previous = mapping.get(key) mapping[key] = entry return self._track_replacement( - kind, key, slot=("manager_mapping", builtins.id(mapping), key), - current=entry, previous=previous, - restore=lambda replacement: self._manager._restore_mapping( - mapping, key, entry, replacement - ), + kind, key, slot=("manager_mapping", id(mapping), key), current=entry, previous=previous, + restore=lambda replacement: self._manager._restore_mapping(mapping, key, entry, replacement), ) def _register_scoped_provider( @@ -418,16 +358,12 @@ class PluginContext: self._llm = PluginLlm(plugin_id=self.plugin_id) return self._llm - @property + @cached_property def subagent_lifecycle(self) -> Any: """Plugin-safe subagent lifecycle service: serializable handles and immutable snapshots, never a live agent or private registry.""" - if self._subagent_lifecycle is None: - from agent.subagent_lifecycle import ( - SubagentLifecycleService, get_active_subagent_parent, - ) - self._subagent_lifecycle = SubagentLifecycleService(get_active_subagent_parent) - return self._subagent_lifecycle + from agent.subagent_lifecycle import SubagentLifecycleService, get_active_subagent_parent + return SubagentLifecycleService(get_active_subagent_parent) @property def profile_name(self) -> str: @@ -470,15 +406,26 @@ class PluginContext: """Register a human approval transport, inactive until ``security.approval.transport: `` selects it. It receives a redacted ``ApprovalRequest`` and returns only a correlated decision; policy and persistence stay host-owned. ``present_fn`` may be async.""" - self._manager.register_approval_transport(name, present_fn, plugin_id=self.plugin_id) - # Record ownership so unload/force-reload removes this transport. Duplicate names are - # rejected above (raise), so there is never a displaced previous entry to restore. + from hermes_cli.approval_transport import RegisteredApprovalTransport + transports = self._manager._approval_transports clean = str(name).strip().lower() - entry = self._manager._approval_transports.get(clean) - if entry is not None: - self._track_mapping_entry( - "approval_transport", clean, self._manager._approval_transports, entry, None - ) + if clean == "builtin": + raise ValueError("approval transport name 'builtin' is reserved") + if not re.fullmatch(r"[a-z0-9][a-z0-9_-]{0,63}", clean): + raise ValueError("approval transport name must match [a-z0-9][a-z0-9_-]{0,63}") + if not callable(present_fn): + raise TypeError("approval transport present_fn must be callable") + if clean in transports: + owner = transports[clean].plugin_id + raise ValueError(f"approval transport {clean!r} is already registered by {owner!r}") + entry = RegisteredApprovalTransport( + name=clean, present=present_fn, plugin_id=self.plugin_id, + profile_home=str(get_hermes_home().resolve()), + ) + logger.info("Plugin %s registered approval transport: %s", self.plugin_id, clean) + # Duplicate names are rejected above, so there is never a displaced previous entry to restore; + # tracking makes unload/force-reload remove this transport. + self._track_mapping_entry("approval_transport", clean, transports, entry, None) @_serialized_replacement def register_tool( @@ -500,8 +447,7 @@ class PluginContext: from tools.registry import registry scope = self._manager.scope_key previous = registry.snapshot_registration(name, scope=scope) - effective = registry.get_entry(name, scope=scope) - if previous is None and effective is not None and not override: + if previous is None and not override and registry.get_entry(name, scope=scope) is not None: logger.warning( "Plugin %s tried to shadow global tool %s without override=True", self.manifest.name, name, @@ -513,16 +459,13 @@ class PluginContext: override=override, scope=scope, ) registered = registry.snapshot_registration(name, scope=scope) - if ( - registered is not None and registered is not previous and registered.handler is handler - ): + handle = None + if registered is not None and registered is not previous and registered.handler is handler: self._manager._plugin_tool_names.add(name) handle = self._manager._track_scoped_registration( self.manifest, "tool", name, registry, registered, previous, finalize=lambda: self._manager._remove_tool_name_if_unowned(name), ) - else: - handle = None logger.debug( "Plugin %s registered tool: %s%s", self.manifest.name, name, " (override)" if override else "", @@ -533,8 +476,7 @@ class PluginContext: """True when *capability* is live for this plugin (probe, then degrade gracefully). Bundled plugins are trusted for ``tools.override``; otherwise granted_capabilities or the legacy ``allow_*`` key decides. Unknown ids / unreadable consent -> False (fail closed).""" - source = getattr(self.manifest, "source", "") or "" - if source == "bundled" and capability == "tools.override": + if self.manifest.source == "bundled" and capability == "tools.override": return True return plugin_capability_granted(self.plugin_id, capability) @@ -543,29 +485,23 @@ class PluginContext: timeout: float = 30, ) -> Dict[str, Any]: """Call ``tool`` on MCP ``server`` synchronously through :mod:`tools.mcp_tool`'s native client - (same trust gates, breaker, reconnect — never a parallel connection). Default-off per-server - grant: servers not in ``plugins.entries..mcp_allowlist`` raise ``PermissionError``. - ``timeout`` clamps to 1–600s. Returns ``{"ok": True, "result"}`` / ``{"ok": False, "error"}``; - results over ~64KB are truncated with a marker.""" - allowlist = self._mcp_allowlist(self.plugin_id) - if server not in allowlist: + (same trust gates, breaker, reconnect — never a parallel connection). Servers not in + ``plugins.entries..mcp_allowlist`` raise ``PermissionError`` (default-deny). ``timeout`` + clamps to 1–600s; results over ~64KB are truncated with a marker.""" + if server not in self._mcp_allowlist(self.plugin_id): raise PermissionError( f"Plugin {self.manifest.name!r} is not allowed to call MCP " f"server {server!r}. Add it to " f"plugins.entries.{self.plugin_id}.mcp_allowlist in config.yaml " f"to grant access (default is no MCP access)." ) - try: timeout = float(timeout) except (TypeError, ValueError): timeout = 30.0 timeout = max(1.0, min(timeout, 600.0)) - from tools.mcp_tool import _make_tool_handler - handler = _make_tool_handler(server, tool, timeout) - raw = handler(dict(arguments or {})) - + raw = _make_tool_handler(server, tool, timeout)(dict(arguments or {})) logger.debug( "Plugin %s called MCP %s/%s (timeout=%ss, %d chars returned)", self.manifest.name, server, tool, timeout, len(raw or ""), @@ -579,14 +515,11 @@ class PluginContext: """Normalize an MCP handler result string into a stable envelope.""" if not isinstance(raw, str): raw = "" if raw is None else str(raw) - if len(raw) > cls._MCP_RESULT_CHAR_CAP: + truncated = len(raw) > cls._MCP_RESULT_CHAR_CAP + if truncated: raw = raw[: cls._MCP_RESULT_CHAR_CAP] + "… [truncated]" - truncated = True - else: - truncated = False - parsed: Any = None try: - parsed = json.loads(raw) + parsed: Any = json.loads(raw) except (ValueError, TypeError): parsed = None if isinstance(parsed, dict) and "error" in parsed: @@ -609,23 +542,17 @@ class PluginContext: cfg = load_config() or {} except Exception: return [] - entries = (cfg.get("plugins") or {}).get("entries") or {} - entry = entries.get(plugin_id) or {} - allowlist = entry.get("mcp_allowlist") - if not isinstance(allowlist, list): - return [] - return [str(item) for item in allowlist] + allowlist = (_plugin_settings_entry(cfg, plugin_id) or {}).get("mcp_allowlist") + return [str(item) for item in allowlist] if isinstance(allowlist, list) else [] def _tool_override_allowed(self, tool_name: str) -> bool: """Whether this plugin may override built-in tools: bundled plugins are trusted (a maintainer choice, not privilege escalation); others need ``tools.override`` via :func:`plugin_capability_granted` (granted_capabilities OR legacy ``allow_tool_override: true``).""" - source = getattr(self.manifest, "source", "") or "" - if source == "bundled": + if self.manifest.source == "bundled": return True try: from hermes_cli.config import load_config - with _plugin_home_scope(self._manager.home_path): cfg = load_config() or {} except Exception: @@ -643,14 +570,10 @@ class PluginContext: request for async dispatch, not that delivery completed.""" cli = self._manager._cli_ref msg = content if role == "user" else f"[{role}] {content}" - if cli is not None: - if getattr(cli, "_agent_running", False): - cli._interrupt_queue.put(msg) - else: - cli._pending_input.put(msg) + queue_ = cli._interrupt_queue if getattr(cli, "_agent_running", False) else cli._pending_input + queue_.put(msg) return True - if not session_key: logger.warning("inject_message: gateway mode requires an existing session_key") return False @@ -661,17 +584,13 @@ class PluginContext: self.plugin_id, ) return False - if not self._manager.has_gateway_message_injector: logger.warning("inject_message: no live gateway is available") return False - try: - return bool( - self._manager.inject_gateway_message( - session_key=session_key, content=msg, plugin_id=self.plugin_id, - ) - ) + return bool(self._manager.inject_gateway_message( + session_key=session_key, content=msg, plugin_id=self.plugin_id, + )) except Exception: logger.warning( "inject_message: gateway scheduling failed for plugin %s", self.plugin_id, @@ -685,12 +604,9 @@ class PluginContext: cfg = load_config_readonly() or {} except Exception: return False - - return ( - cfg_get( - cfg, "plugins", "entries", self.plugin_id, "allow_gateway_injection", default=False, - ) is True - ) + return (_plugin_settings_entry(cfg, self.plugin_id) or {}).get( + "allow_gateway_injection" + ) is True @_serialized_replacement def register_cli_command( @@ -699,14 +615,11 @@ class PluginContext: ) -> PluginRegistration: """Register a CLI subcommand (``hermes ...``). *setup_fn* receives the argparse subparser; *handler_fn* becomes ``set_defaults(func=...)``.""" - previous = self._manager._cli_commands.get(name) entry = { "name": name, "help": help, "description": description, "setup_fn": setup_fn, "handler_fn": handler_fn, "plugin": self.manifest.name, "plugin_key": self.plugin_id, } - handle = self._track_mapping_entry( - "cli_command", name, self._manager._cli_commands, entry, previous - ) + handle = self._track_mapping_entry("cli_command", name, self._manager._cli_commands, entry) logger.debug("Plugin %s registered CLI command: %s", self.manifest.name, name) return handle @@ -715,19 +628,16 @@ class PluginContext: self, name: str, handler: Callable, description: str = "", args_hint: str = "", argument_mode: str | None = None, ) -> Optional[PluginRegistration]: - """Register an in-session slash command (``/name``) for CLI and gateway sessions. Handler: - ``fn(raw_args: str) -> str | None`` (sync or async). ``args_hint`` (e.g. ``""``) lets - adapters like Discord surface an argument field; without it the command registers parameterless - there but still accepts trailing text as free-form chat.""" + """Register an in-session slash command (``/name``); handler ``fn(raw_args: str) -> str | None`` + (sync or async). ``args_hint`` (e.g. ``""``) lets adapters like Discord surface an argument + field; without it the command registers parameterless there but still accepts trailing text.""" clean = name.lower().strip().lstrip("/").replace(" ", "-") if not clean: logger.warning( "Plugin '%s' tried to register a command with an empty name.", self.manifest.name, ) return - - # Reject if it conflicts with a built-in command - try: + with suppress(Exception): # reject if it conflicts with a built-in command from hermes_cli.commands import resolve_command if resolve_command(clean) is not None: logger.warning( @@ -735,10 +645,6 @@ class PluginContext: "with a built-in command. Skipping.", self.manifest.name, clean, ) return - except Exception: - pass - - previous = self._manager._plugin_commands.get(clean) hint = (args_hint or "").strip() mode = argument_mode if argument_mode in {"options", "text", "mixed"} else ( "text" if hint else None @@ -748,9 +654,7 @@ class PluginContext: "plugin": self.manifest.name, "plugin_key": self.plugin_id, "args_hint": hint, "argument_mode": mode, } - handle = self._track_mapping_entry( - "command", clean, self._manager._plugin_commands, entry, previous - ) + handle = self._track_mapping_entry("command", clean, self._manager._plugin_commands, entry) logger.debug("Plugin %s registered command: /%s", self.manifest.name, clean) return handle @@ -760,11 +664,9 @@ class PluginContext: from tools.registry import registry # In gateway mode _cli_ref is None — tools degrade gracefully (no spinner, TERMINAL_CWD). if "parent_agent" not in kwargs: - cli = self._manager._cli_ref - agent = getattr(cli, "agent", None) if cli else None + agent = getattr(self._manager._cli_ref, "agent", None) if agent is not None: kwargs["parent_agent"] = agent - return registry.dispatch(tool_name, args, scope=self._manager.scope_key, **kwargs) @_serialized_replacement @@ -868,10 +770,9 @@ class PluginContext: ) -> Optional[PluginRegistration]: """Register a gateway platform adapter (``adapter_factory(PlatformConfig) -> BasePlatformAdapter``). ``check_fn`` is a PASSIVE "deps importable?" probe that must never install (status displays call - it freely); pass an ACTIVE installer as ``ensure_deps_fn`` (the gateway calls it from - ``create_adapter()`` when ``check_fn`` is False). Extra kwargs (``setup_fn``, ``emoji``, - ``allowed_users_env``, ``platform_hint``, ``ensure_deps_fn``) forward to ``PlatformEntry``; - unknown keys raise TypeError.""" + it freely); an ACTIVE installer goes in ``ensure_deps_fn`` (called from ``create_adapter()`` when + ``check_fn`` is False). Extra kwargs (``setup_fn``, ``emoji``, ``allowed_users_env``, + ``platform_hint``, ``ensure_deps_fn``) forward to ``PlatformEntry``; unknown keys raise TypeError.""" from gateway.platform_registry import platform_registry, PlatformEntry entry_kwargs.setdefault("plugin_name", self.manifest.name) entry = PlatformEntry( @@ -905,24 +806,23 @@ class PluginContext: if action_id is None or (isinstance(action_id, str) and not action_id.strip()): raise self._refuse("a Slack action handler with an empty action_id") entry = (action_id, callback, self.manifest.name) - self._manager._slack_action_handlers.append(entry) + handlers = self._manager._slack_action_handlers + handlers.append(entry) handle = self._track( "slack_action_handler", repr(action_id), - lambda: self._manager._remove_identity( - self._manager._slack_action_handlers, entry - ), + lambda: self._manager._remove_identity(handlers, entry), ) logger.debug("Plugin %s registered Slack action handler: %s", self.manifest.name, action_id) return handle def register_platform_handler(self, platform: str, factory: Callable) -> None: - """Register a native-client handler factory for a gateway platform, invoked at ``connect()`` as - ``factory(native, adapter)`` before/as the core handlers register (``adapter`` read-only). - ``native``: telegram PTB ``Application``, discord ``commands.Bot``, slack ``AsyncApp``, matrix - client, teams ``App``, dingtalk ``DingTalkStreamClient``, line aiohttp ``web.Application``, - others ``None``. Keep SDK imports inside the factory; exceptions are logged and the platform - still connects. Always scope handlers hooked into first-match dispatch tables so core flows - keep working. Raises ``ValueError`` for a non-callable factory or empty platform.""" + """Register ``factory(native, adapter)``, invoked at ``connect()`` before/as the core handlers + register (``adapter`` read-only). ``native``: telegram PTB ``Application``, discord + ``commands.Bot``, slack ``AsyncApp``, matrix client, teams ``App``, dingtalk + ``DingTalkStreamClient``, line aiohttp ``web.Application``, others ``None``. Keep SDK imports + inside the factory; exceptions are logged and the platform still connects. Scope handlers in + first-match dispatch tables so core flows keep working. Raises ``ValueError`` when not callable + or platform is empty.""" if not callable(factory): raise self._refuse("a platform handler factory with a non-callable factory") key = (platform or "").strip().lower() @@ -937,10 +837,9 @@ class PluginContext: ) def register_telegram_handler(self, factory: Callable) -> None: - """Alias of ``register_platform_handler("telegram", factory)``: ``factory(application, - adapter)`` runs before the core handlers. PTB dispatches only the FIRST matching handler per - group and core registers a catch-all ``CallbackQueryHandler`` — always scope with - ``pattern=`` or you swallow the core button flows. Raises ``ValueError`` if not callable.""" + """``register_platform_handler("telegram", factory)``. PTB dispatches only the FIRST matching + handler per group and core registers a catch-all ``CallbackQueryHandler`` — always scope with + ``pattern=`` or you swallow the core button flows.""" self.register_platform_handler("telegram", factory) @_serialized_replacement @@ -949,10 +848,9 @@ class PluginContext: defaults: Optional[Dict[str, Any]] = None, ) -> PluginRegistration: """Register an auxiliary LLM task with its own ``auxiliary.`` config block (picker entry, - ``AUXILIARY__*`` env bridge, defaults merged into loaded configs). ``key`` is snake_case - and must not shadow a built-in task; ``defaults`` may override - provider/model/base_url/api_key/timeout/extra_body (unknown keys preserved verbatim). Raises - ``ValueError`` for an empty/invalid key, a built-in key, or another plugin's key.""" + ``AUXILIARY__*`` env bridge, defaults merged into loaded configs). ``defaults`` may + override provider/model/base_url/api_key/timeout/extra_body (unknown keys kept verbatim). + Raises ``ValueError`` for an empty/invalid key, a built-in key, or another plugin's key.""" if not key or not isinstance(key, str): raise ValueError( f"Plugin '{self.manifest.name}' tried to register auxiliary task with invalid key {key!r}" @@ -962,33 +860,26 @@ class PluginContext: f"Plugin '{self.manifest.name}' auxiliary task key {key!r} " f"must contain only alphanumeric characters and underscores" ) - from hermes_cli.main import _AUX_TASKS as _BUILTIN_AUX_TASKS - builtin_keys = {k for k, _name, _desc in _BUILTIN_AUX_TASKS} - if key in builtin_keys: + if key in {k for k, _name, _desc in _BUILTIN_AUX_TASKS}: raise ValueError( f"Plugin '{self.manifest.name}' cannot register auxiliary task " f"{key!r} — that key is reserved for a built-in task. " f"Pick a plugin-namespaced key (e.g. '{self.manifest.name}_{key}')." ) - # Owner is the canonical id ``ctx.llm`` is bound to, so agent/plugin_llm.py can match it. owner_id = self.plugin_id - existing = self._manager._aux_tasks.get(key) if existing is not None and existing.get("plugin") != owner_id: raise ValueError( f"Plugin '{self.manifest.name}' cannot register auxiliary task " f"{key!r} — already registered by plugin " f"'{existing.get('plugin')}'" ) - # Plugin owns the schema; routing fields are guaranteed present so consumers don't crash. merged_defaults: Dict[str, Any] = { "provider": "auto", "model": "", "base_url": "", "api_key": "", - "timeout": 60, "extra_body": {}, + "timeout": 60, "extra_body": {}, **(defaults or {}), } - merged_defaults.update(defaults or {}) - entry = { "key": key, "display_name": display_name, "description": description, "defaults": merged_defaults, "plugin": owner_id, "plugin_key": owner_id, @@ -1021,6 +912,13 @@ class PluginContext: """Register a lifecycle hook callback (unknown names warn but are still stored).""" return self._track_callback("hook", hook_name, callback, self._manager._hooks, VALID_HOOKS) + def register_middleware(self, kind: str, callback: Callable) -> PluginRegistration: + """Register behavior-changing middleware (request kinds rewrite the payload, execution kinds + wrap the callback). Unknown kinds warn but are stored.""" + return self._track_callback( + "middleware", kind, callback, self._manager._middleware, VALID_MIDDLEWARE + ) + def _track_callback( self, kind: str, key: str, callback: Callable, mapping: Dict[str, List[Callable]], valid: Set[str], @@ -1102,8 +1000,7 @@ class PluginContext: ) 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 {}) + return self._manager._dispatch_event(f"{plugin_key}:{event}", payload or {}) def subscribe(self, event: str, callback: Callable) -> None: """Subscribe to a fully-qualified ``:`` name (unrestricted — only @@ -1112,17 +1009,9 @@ class PluginContext: raise ValueError( f"Plugin '{self.manifest.name}' subscribe() requires a " f"non-empty event name" ) - plugin_key = self.plugin_id - self._manager._subscribe_event(plugin_key, event, callback) + self._manager._subscribe_event(self.plugin_id, event, callback) logger.debug("Plugin %s subscribed to event: %s", self.manifest.name, event) - def register_middleware(self, kind: str, callback: Callable) -> PluginRegistration: - """Register behavior-changing middleware (request kinds rewrite the payload, execution kinds - wrap the callback). Unknown kinds warn but are stored.""" - return self._track_callback( - "middleware", kind, callback, self._manager._middleware, VALID_MIDDLEWARE - ) - @_serialized_replacement def register_skill( self, name: str, path: Path, description: str = "", @@ -1142,19 +1031,15 @@ class PluginContext: raise ValueError(f"Invalid skill name '{name}'. Must match [a-zA-Z0-9_-]+.") if not path.exists(): raise FileNotFoundError(f"SKILL.md not found at {path}") - namespace = self.manifest.skill_namespace or self.manifest.name qualified = f"{namespace}:{name}" if self.manifest.portable and qualified in self._manager._plugin_skills: raise ValueError(f"Plugin skill '{qualified}' is already registered") - previous = self._manager._plugin_skills.get(qualified) entry = { "path": path, "plugin": namespace, "plugin_key": self.plugin_id, "bare_name": name, "description": description, "frontmatter": dict(frontmatter or {}), } - handle = self._track_mapping_entry( - "skill", qualified, self._manager._plugin_skills, entry, previous - ) + handle = self._track_mapping_entry("skill", qualified, self._manager._plugin_skills, entry) logger.debug("Plugin %s registered skill: %s", self.manifest.name, qualified) return handle @@ -1162,91 +1047,83 @@ class PluginContext: # -- scoped provider registrars ------------------------------------------------------------------ # Every ``register__provider`` shares one body (:meth:`PluginContext._register_scoped_provider`): # type-check, register in the scope-keyed process-global registry, lease the slot so unload restores -# the displaced entry. Rows: (method, kind, registry module, base-class module:attr, label, options). -# ``normalize`` defaults to ``str.strip``; ``None`` keeps the raw name; ``lower`` also lowercases. -_SCOPED_PROVIDER_REGISTRARS: Tuple[Tuple[str, str, str, str, str, Dict[str, Any]], ...] = ( +# the displaced entry. Rows: (method, kind, registry module, base-class module:attr, label, docstring, +# options). ``normalize``: ``strip`` (default), ``lower`` (strip+lowercase) or ``None`` (raw name). +_SCOPED_PROVIDER_REGISTRARS: Tuple[Tuple[str, str, str, str, str, str, Dict[str, Any]], ...] = ( ("register_image_gen_provider", "image_gen_provider", "agent.image_gen_registry", - "agent.image_gen_provider:ImageGenProvider", "image_gen provider", {"article": "an"}), + "agent.image_gen_provider:ImageGenProvider", "image_gen provider", + "Register an :class:`agent.image_gen_provider.ImageGenProvider`; " + "``provider.name`` is matched by ``image_gen.provider``.", {"article": "an"}), ("register_video_gen_provider", "video_gen_provider", "agent.video_gen_registry", - "agent.video_gen_provider:VideoGenProvider", "video_gen provider", {}), + "agent.video_gen_provider:VideoGenProvider", "video_gen provider", + "Register an :class:`agent.video_gen_provider.VideoGenProvider`; " + "``provider.name`` is matched by ``video_gen.provider``.", {}), ("register_web_search_provider", "web_search_provider", "agent.web_search_registry", - "agent.web_search_provider:WebSearchProvider", "web provider", {}), + "agent.web_search_provider:WebSearchProvider", "web provider", + "Register an :class:`agent.web_search_provider.WebSearchProvider`; " + "``provider.name`` is matched by ``web.search_backend`` / ``web.extract_backend`` / ``web.backend``.", + {}), ("register_browser_provider", "browser_provider", "agent.browser_registry", - "agent.browser_provider:BrowserProvider", "browser provider", {}), + "agent.browser_provider:BrowserProvider", "browser provider", + "Register an :class:`agent.browser_provider.BrowserProvider`; " + "``provider.name`` is matched by ``browser.cloud_provider`` (consulted by " + "``tools.browser_tool._get_cloud_provider``).", {}), ("register_terminal_environment_provider", "terminal_environment_provider", "agent.terminal_env_registry", "agent.terminal_env_provider:TerminalEnvironmentProvider", "terminal environment provider", + "Register a :class:`agent.terminal_env_provider.TerminalEnvironmentProvider`; ``provider.name`` " + "is matched by ``terminal.backend`` when no built-in backend has that name. Built-in names (local, " + "docker, singularity, modal, daytona, vercel_sandbox, ssh) are rejected — plugins never shadow " + "in-tree backends.", {"normalize": "lower", "reject_message": "Plugin '%s' terminal environment provider rejected: %s"}), ("register_secret_source", "secret_source", "agent.secret_sources.registry", "agent.secret_sources.base:SecretSource", "secret source", + "Register a :class:`agent.secret_sources.base.SecretSource`, run by ``load_hermes_dotenv()`` " + "(after ``~/.hermes/.env``, before credentials are read) when ``secrets.`` is enabled. The " + "orchestrator owns ordering/precedence/provenance; the source only fetches. Since dotenv usually " + "loads before discovery, the manager re-pulls enabled plugin sources afterwards.", {"normalize": None, "register": "register_source", "param": "source"}), ("register_tts_provider", "tts_provider", "agent.tts_registry", - "agent.tts_provider:TTSProvider", "TTS provider", {"normalize": "lower"}), + "agent.tts_provider:TTSProvider", "TTS provider", + "Register an :class:`agent.tts_provider.TTSProvider`; ``provider.name`` is matched by " + "``tts.provider`` unless it is a built-in name (rejected with a warning) or a " + "``tts.providers.: type: command`` entry shares it (command-providers win).", + {"normalize": "lower"}), ("register_transcription_provider", "transcription_provider", "agent.transcription_registry", "agent.transcription_provider:TranscriptionProvider", "transcription provider", - {"normalize": "lower"}), + "Register an :class:`agent.transcription_provider.TranscriptionProvider`; ``provider.name`` is " + "matched by ``stt.provider`` unless it is a built-in name (rejected) or a ``stt.providers.: " + "type: command`` entry shares it (command-providers win).", {"normalize": "lower"}), ) -_SCOPED_PROVIDER_DOCS: Dict[str, str] = { - "register_image_gen_provider": "Register an :class:`agent.image_gen_provider.ImageGenProvider`; " - "``provider.name`` is matched by ``image_gen.provider``.", - "register_video_gen_provider": "Register an :class:`agent.video_gen_provider.VideoGenProvider`; " - "``provider.name`` is matched by ``video_gen.provider``.", - "register_web_search_provider": "Register an :class:`agent.web_search_provider.WebSearchProvider`; " - "``provider.name`` is matched by ``web.search_backend`` / ``web.extract_backend`` / ``web.backend``.", - "register_browser_provider": "Register an :class:`agent.browser_provider.BrowserProvider`; " - "``provider.name`` is matched by ``browser.cloud_provider`` (consulted by " - "``tools.browser_tool._get_cloud_provider``).", - "register_terminal_environment_provider": "Register a " - ":class:`agent.terminal_env_provider.TerminalEnvironmentProvider`; ``provider.name`` is matched " - "by ``terminal.backend`` when no built-in backend has that name. Built-in names (local, docker, " - "singularity, modal, daytona, vercel_sandbox, ssh) are rejected — plugins never shadow in-tree " - "backends.", - "register_secret_source": "Register a :class:`agent.secret_sources.base.SecretSource`, run by " - "``load_hermes_dotenv()`` (after ``~/.hermes/.env``, before credentials are read) when " - "``secrets.`` is enabled. The orchestrator owns ordering/precedence/provenance; the source " - "only fetches. Since dotenv usually loads before discovery, the manager re-pulls enabled plugin " - "sources afterwards.", - "register_tts_provider": "Register an :class:`agent.tts_provider.TTSProvider`; ``provider.name`` is " - "matched by ``tts.provider`` unless it is a built-in name (rejected with a warning) or a " - "``tts.providers.: type: command`` entry shares it (command-providers win).", - "register_transcription_provider": "Register an " - ":class:`agent.transcription_provider.TranscriptionProvider`; ``provider.name`` is matched by " - "``stt.provider`` unless it is a built-in name (rejected) or a ``stt.providers.: type: " - "command`` entry shares it (command-providers win).", +_NAME_NORMALIZERS: Dict[Optional[str], Optional[Callable[[str], str]]] = { + "strip": lambda n: n.strip(), "lower": lambda n: n.strip().lower(), None: None, } -def _make_scoped_provider_registrar(method_name, kind, registry_mod, base_ref, label, options): +def _make_scoped_provider_registrar(method_name, kind, registry_mod, base_ref, label, doc, options): """Build one ``register__provider`` method from a ``_SCOPED_PROVIDER_REGISTRARS`` row.""" base_mod, base_attr = base_ref.split(":") - normalize = options.get("normalize", "strip") - normalize_fn = ( - None if normalize is None else (lambda n: n.strip().lower()) if normalize == "lower" - else (lambda n: n.strip()) - ) + normalize_fn = _NAME_NORMALIZERS[options.get("normalize", "strip")] + register_name = options.get("register") def register(self, provider) -> Optional[PluginRegistration]: - return _register(self, provider) - - def register_source(self, source) -> Optional[PluginRegistration]: # secret sources: ``source`` - return _register(self, source) - - def _register(self, provider): registry = importlib.import_module(registry_mod) - base_class = getattr(importlib.import_module(base_mod), base_attr) - register_fn = options.get("register") return self._register_scoped_provider( - provider, kind=kind, base_class=base_class, registry=registry, label=label, - article=options.get("article", "a"), normalize=normalize_fn, - register=getattr(registry, register_fn) if register_fn else None, + provider, kind=kind, base_class=getattr(importlib.import_module(base_mod), base_attr), + registry=registry, label=label, article=options.get("article", "a"), + normalize=normalize_fn, + register=getattr(registry, register_name) if register_name else None, reject_message=options.get("reject_message"), ) + def register_source(self, source) -> Optional[PluginRegistration]: # secret sources: ``source`` + return register(self, source) + method = register_source if options.get("param") == "source" else register method.__name__ = method_name method.__qualname__ = f"PluginContext.{method_name}" - method.__doc__ = _SCOPED_PROVIDER_DOCS[method_name] + method.__doc__ = doc return _serialized_replacement(method) @@ -1261,21 +1138,17 @@ def _resolve_hook_callback_timeout() -> float: timeout = _HOOK_CALLBACK_TIMEOUT_SECS try: from hermes_cli.config import load_config_readonly - plugins_cfg = (load_config_readonly() or {}).get("plugins") - if isinstance(plugins_cfg, dict) and "hook_callback_timeout" in plugins_cfg: - raw = plugins_cfg.get("hook_callback_timeout") - if raw is not None: - timeout = float(raw) + if isinstance(plugins_cfg, dict) and plugins_cfg.get("hook_callback_timeout") is not None: + timeout = float(plugins_cfg["hook_callback_timeout"]) except (TypeError, ValueError): logger.warning( "plugins.hook_callback_timeout is not a number; using default %gs", _HOOK_CALLBACK_TIMEOUT_SECS, ) - timeout = _HOOK_CALLBACK_TIMEOUT_SECS + return _HOOK_CALLBACK_TIMEOUT_SECS except Exception: - timeout = _HOOK_CALLBACK_TIMEOUT_SECS - + return _HOOK_CALLBACK_TIMEOUT_SECS if timeout < 0: logger.warning( "plugins.hook_callback_timeout=%g is negative; using default %gs", timeout, @@ -1300,22 +1173,25 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): self.scope_key = scope_key or hermes_home_key() self.home_path = Path(self.scope_key) self._discovery_lock = threading.RLock() + self._discovered: bool = False + self._cli_ref = None # Set by CLI after plugin discovery + self._gateway_message_injector: tuple[object, Callable] | None = None + self._context_engine = None # Set by a plugin via register_context_engine() + # Manager-local registries keyed by name (see the matching ``PluginContext.register_*``). self._plugins: Dict[str, LoadedPlugin] = {} self._hooks: Dict[str, List[Callable]] = {} self._middleware: Dict[str, List[Callable]] = {} self._plugin_tool_names: Set[str] = set() self._plugin_platform_names: Set[str] = set() self._cli_commands: Dict[str, dict] = {} - self._context_engine = None # Set by a plugin via register_context_engine() - self._plugin_commands: Dict[str, dict] = {} # Slash commands registered by plugins + self._plugin_commands: Dict[str, dict] = {} # slash commands self._system_prompt_sections: Dict[str, PluginSystemPromptSection] = {} - self._discovered: bool = False - self._cli_ref = None # Set by CLI after plugin discovery - self._gateway_message_injector: tuple[object, Callable] | None = None self._plugin_skills: Dict[str, Dict[str, Any]] = {} # qualified name -> metadata self._portable_mcp_servers: Dict[str, Dict[str, Any]] = {} - self._aux_tasks: Dict[str, Dict[str, Any]] = {} # see register_auxiliary_task + self._aux_tasks: Dict[str, Dict[str, Any]] = {} self._approval_transports: Dict[str, Any] = {} + self._slack_action_handlers: List[tuple] = [] # (matcher, callback, plugin_name) + self._platform_handler_factories: Dict[str, List[tuple]] = {} # lowercase platform -> list # Event bus: owner-tagged subscriptions (unload removes zombies); one daemon worker keeps # registration order while emitters never block. self._subscriptions: Dict[str, List[_EventSubscription]] = {} @@ -1326,7 +1202,6 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): self._event_queue: queue.Queue[Any] = queue.Queue(maxsize=_EVENT_PENDING_CAP) self._event_worker: Optional[threading.Thread] = None self._emit_depth = threading.local() # per-worker chain depth caps mutual emitters - self._slack_action_handlers: List[tuple] = [] # (matcher, callback, plugin_name) # In-flight / recently-timed-out hook callbacks keyed by (hook_name, id(cb)) so a stuck # policy hook cannot spawn a new abandoned thread on every fire. self._hook_running_callbacks: Dict[tuple, object] = {} @@ -1347,8 +1222,6 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): # and contributed tool names (so `hermes plugins list` still attributes them). self._predeclared_modules: Dict[str, types.ModuleType] = {} self._predeclared_tools: Dict[str, List[str]] = {} - # Native platform handler factories keyed by lowercase platform name. - self._platform_handler_factories: Dict[str, List[tuple]] = {} @property def has_gateway_message_injector(self) -> bool: @@ -1368,9 +1241,7 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): def inject_gateway_message(self, **kwargs: Any) -> bool: """Submit a plugin-triggered turn to the live gateway.""" registered = self._gateway_message_injector - if registered is None: - return False - return bool(registered[1](**kwargs)) + return registered is not None and bool(registered[1](**kwargs)) def discover_and_load(self, force: bool = False) -> None: """Scan all plugin sources and load each plugin found; ``force`` unloads first so config @@ -1406,20 +1277,12 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): def _re_register_config_hooks_after_force(self) -> None: """Restore config-owned shell hooks/outbound webhooks after a force clear; each guarded independently so one failing does not skip the other.""" - try: - from agent.shell_hooks import re_register_config_hooks - - re_register_config_hooks() - except Exception as exc: - logger.debug("force-reload shell-hook re-register skipped: %s", exc) - try: - from agent.outbound_webhooks import ( - re_register_config_hooks as re_register_outbound_webhooks, - ) - - re_register_outbound_webhooks() - except Exception as exc: - logger.debug("force-reload outbound-webhook re-register skipped: %s", exc) + for label, module_name in (("shell-hook", "agent.shell_hooks"), + ("outbound-webhook", "agent.outbound_webhooks")): + try: + importlib.import_module(module_name).re_register_config_hooks() + except Exception as exc: + logger.debug("force-reload %s re-register skipped: %s", label, exc) def _refresh_secret_sources_after_discovery(self) -> None: """If any plugin secret source is enabled (per its own ``is_enabled(cfg)``, honoring custom @@ -1427,9 +1290,6 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): try: from agent.secret_sources.registry import list_plugin_sources from hermes_cli.env_loader import load_hermes_dotenv, reset_secret_source_cache - except Exception: - return - try: plugin_sources = list_plugin_sources() except Exception: return @@ -1437,18 +1297,15 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): return try: from hermes_cli.config import load_config - - cfg = load_config() or {} - secrets = cfg.get("secrets") or {} + secrets = (load_config() or {}).get("secrets") or {} except Exception: secrets = {} enabled_names = [] for source in plugin_sources: name = getattr(source, "name", "") section = secrets.get(name) - section = section if isinstance(section, dict) else {} try: - if source.is_enabled(section): + if source.is_enabled(section if isinstance(section, dict) else {}): enabled_names.append(name) except Exception: continue # mirrors the orchestrator: a raising is_enabled() is skipped @@ -1472,7 +1329,6 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): ep_manifests = self._scan_entry_points() logger.debug(" entrypoints: %d manifest(s)", len(ep_manifests)) manifests.extend(ep_manifests) - disabled = _get_disabled_plugins() enabled = _get_enabled_plugins() # None = opt-in default (nothing enabled) stale_relay_keys = legacy_relay_plugin_keys(enabled) @@ -1494,27 +1350,18 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): self._warn_python_dependencies(manifest) self._validate_plugin_config_schema(manifest) self._load_plugin(manifest) - if manifests: logger.info( "Plugin discovery complete: %d found, %d enabled", len(self._plugins), sum(1 for p in self._plugins.values() if p.enabled), ) - def _record_placeholder( - self, manifest: PluginManifest, *, enabled: bool, error: Optional[str] = None - ) -> None: - """Record a manifest that discovery will not import (introspection-only entry).""" - loaded = LoadedPlugin(manifest=manifest, enabled=enabled) - if error is not None: - loaded.error = error - self._plugins[manifest_key(manifest)] = loaded - def _gate_manifest( self, manifest: PluginManifest, disabled: Set[str], enabled: Optional[Set[str]], ) -> bool: """Route one winning manifest per :func:`gate_manifest`: load now, defer, or record as - skipped. Returns True only for plugins that go through the dependency-ordered load pass.""" + skipped (introspection-only placeholder). Returns True only for plugins that go through the + dependency-ordered load pass.""" verdict = gate_manifest(manifest, disabled, enabled) if verdict.action == "load": return True @@ -1523,38 +1370,17 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): elif verdict.action == "defer": self._register_deferred_platform(manifest) else: - self._record_placeholder(manifest, enabled=verdict.enabled, error=verdict.error) + self._plugins[manifest_key(manifest)] = LoadedPlugin( + manifest=manifest, enabled=verdict.enabled, error=verdict.error, + ) if verdict.log: logger.log(*verdict.log) return False - def register_approval_transport( - self, name: str, present_fn: Callable, *, plugin_id: str, - ) -> None: - """Register one plugin-owned approval transport for this profile.""" - from hermes_cli.approval_transport import RegisteredApprovalTransport - clean = str(name).strip().lower() - if clean == "builtin": - raise ValueError("approval transport name 'builtin' is reserved") - if not re.fullmatch(r"[a-z0-9][a-z0-9_-]{0,63}", clean): - raise ValueError("approval transport name must match [a-z0-9][a-z0-9_-]{0,63}") - if not callable(present_fn): - raise TypeError("approval transport present_fn must be callable") - if clean in self._approval_transports: - owner = self._approval_transports[clean].plugin_id - raise ValueError(f"approval transport {clean!r} is already registered by {owner!r}") - self._approval_transports[clean] = RegisteredApprovalTransport( - name=clean, present=present_fn, plugin_id=plugin_id, - profile_home=str(get_hermes_home().resolve()), - ) - logger.info("Plugin %s registered approval transport: %s", plugin_id, clean) - def get_approval_transport(self, name: str): """Return a transport only inside the profile that registered it.""" registered = self._approval_transports.get(str(name).strip().lower()) - if registered is None: - return None - if registered.profile_home != str(get_hermes_home().resolve()): + if registered is None or registered.profile_home != str(get_hermes_home().resolve()): return None return registered @@ -1567,7 +1393,6 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): manifest collection so precedence/gating cannot diverge).""" if _env_enabled("HERMES_SAFE_MODE"): return False - plugins_config = raw_config.get("plugins") if not isinstance(plugins_config, dict): return False @@ -1579,19 +1404,13 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): disabled = _names(plugins_config.get("disabled", [])) if not enabled: return False - winners = {manifest_key(m): m for m in self._collect_directory_manifests()} - for manifest in winners.values(): - if not manifest.portable: - continue - lookup_key = manifest_key(manifest) - if lookup_key in disabled or manifest.name in disabled: - continue - if lookup_key not in enabled and manifest.name not in enabled: + for lookup_key, manifest in winners.items(): + names = {lookup_key, manifest.name} + if not manifest.portable or names & disabled or not names & enabled: continue try: from hermes_cli.agent_plugins import _discover_mcp - if _discover_mcp( Path(manifest.path), get_hermes_home() / "plugin-data" / (manifest.skill_namespace or lookup_key), [], create_data=False, @@ -1618,25 +1437,21 @@ class PluginManager(PluginLoaderMixin, PluginDispatchMixin, PluginLedgerMixin): def get_platform_handler_factories(self, platform: str) -> List[tuple]: """``(factory, plugin_name)`` tuples for one platform; adapters call ``factory(native, adapter)`` at connect (see :meth:`PluginContext.register_platform_handler`).""" - key = (platform or "").strip().lower() - return list(self._platform_handler_factories.get(key, [])) + return list(self._platform_handler_factories.get((platform or "").strip().lower(), [])) def list_plugins(self) -> List[Dict[str, Any]]: """Return a list of info dicts for all discovered plugins.""" - result: List[Dict[str, Any]] = [] - for key, loaded in sorted(self._plugins.items()): - result.append( - { - "name": loaded.manifest.name, "key": manifest_key(loaded.manifest), - "kind": loaded.manifest.kind, "version": loaded.manifest.version, - "description": loaded.manifest.description, "source": loaded.manifest.source, - "enabled": loaded.enabled, "tools": len(loaded.tools_registered), - "hooks": len(loaded.hooks_registered), - "middleware": len(loaded.middleware_registered), - "commands": len(loaded.commands_registered), "error": loaded.error, - } - ) - return result + return [ + { + "name": loaded.manifest.name, "key": manifest_key(loaded.manifest), + "kind": loaded.manifest.kind, "version": loaded.manifest.version, + "description": loaded.manifest.description, "source": loaded.manifest.source, + "enabled": loaded.enabled, "tools": len(loaded.tools_registered), + "hooks": len(loaded.hooks_registered), + "middleware": len(loaded.middleware_registered), + "commands": len(loaded.commands_registered), "error": loaded.error, + } for _key, loaded in sorted(self._plugins.items()) + ] def find_plugin_skill(self, qualified_name: str) -> Optional[Path]: """Return the ``Path`` to a plugin skill's SKILL.md, or ``None``.""" @@ -1697,8 +1512,7 @@ def _clear_plugin_submodules(manager: Optional[PluginManager]) -> None: if manager is None: return for loaded in getattr(manager, "_plugins", {}).values(): - module = getattr(loaded, "module", None) - module_name = getattr(module, "__name__", None) + module_name = getattr(getattr(loaded, "module", None), "__name__", None) if not module_name or not module_name.startswith(f"{_NS_PARENT}."): continue _evict_modules(module_name) @@ -1712,7 +1526,6 @@ def get_plugin_manager() -> PluginManager: profile switch gets its own manager and plugin submodules).""" global _plugin_manager current_home = _plugin_home_key() - with _plugin_managers_lock: # Tests/embedders monkeypatch ``_plugin_manager`` directly: adopt a single-slot manager the # keyed cache doesn't know about at all. @@ -1721,12 +1534,10 @@ def get_plugin_manager() -> PluginManager: ): _plugin_managers_by_home[current_home] = _plugin_manager return _plugin_manager - manager = _plugin_managers_by_home.get(current_home) if manager is None: manager = PluginManager(scope_key=hermes_home_key(current_home)) _plugin_managers_by_home[current_home] = manager - _plugin_manager = manager return manager @@ -1749,11 +1560,8 @@ def _reset_plugin_managers_for_tests() -> None: # Dashboard-auth providers are persistent and survive a routine unload, so the clean-slate # reset must clear that process-global registry explicitly or a test's provider leaks. try: - from hermes_cli.dashboard_auth.registry import ( - clear_providers as _clear_dashboard_auth_providers, - ) - - _clear_dashboard_auth_providers() + from hermes_cli.dashboard_auth.registry import clear_providers + clear_providers() except Exception: logger.debug("dashboard-auth registry clear failed", exc_info=True) @@ -1808,8 +1616,7 @@ def _join_background_discovery(timeout: float = 30.0) -> None: t.join(timeout=timeout) -def _plugin_toolset_keys_cache_path(): - from hermes_constants import get_hermes_home +def _plugin_toolset_keys_cache_path() -> Path: return get_hermes_home() / "cache" / "plugin_toolset_keys.json" @@ -1832,26 +1639,17 @@ def _persist_plugin_toolset_keys() -> None: logger.debug("plugin toolset key persist failed", exc_info=True) -def _read_plugin_keys_cache() -> Optional[dict]: - try: - blob = json.loads(_plugin_toolset_keys_cache_path().read_text(encoding="utf-8")) - if isinstance(blob, dict): - return blob - except Exception: - pass - return None - - def _nowait_plugin_set(cache_field: str, live: Callable[[PluginManager], "set[str]"]) -> "set[str]": """Shared body of the ``*_nowait`` probes: live registry, else last launch's cache, else block.""" manager = get_plugin_manager() t = _background_discovery_thread - if manager._discovered and (t is None or not t.is_alive()): + in_flight = t is not None and t.is_alive() + if manager._discovered and not in_flight: return live(manager) - if t is not None and t.is_alive(): - blob = _read_plugin_keys_cache() - if blob is not None: - values = blob.get(cache_field) + if in_flight: + with suppress(Exception): + blob = json.loads(_plugin_toolset_keys_cache_path().read_text(encoding="utf-8")) + values = blob.get(cache_field) if isinstance(blob, dict) else None if isinstance(values, list) and all(isinstance(v, str) for v in values): return set(values) discover_plugins() @@ -1903,9 +1701,7 @@ def invoke_middleware(kind: str, **kwargs: Any) -> List[Any]: def has_middleware(kind: str) -> bool: """True when middleware is registered for ``kind``; lazy-discovers first since callers gate :func:`invoke_middleware` on it.""" - manager = get_plugin_manager() - if not getattr(manager, "_discovered", True): - manager = _delivery_manager() + manager = _delivery_manager() method = getattr(manager, "has_middleware", None) if callable(method): return bool(method(kind)) @@ -1983,16 +1779,13 @@ def _get_pre_tool_call_directive_details( if allowed is not None and tool_name not in allowed: fmt = getattr(_thread_tool_whitelist, "fmt", "Tool '{tool_name}' denied") return _PreToolCallDirective(action="block", message=fmt.format(tool_name=tool_name)) - from hermes_cli.lifecycle import invoke_hook as invoke_lifecycle_hook hook_results = invoke_lifecycle_hook( "pre_tool_call", tool_name=tool_name, args=args if isinstance(args, dict) else {}, task_id=task_id, session_id=session_id, tool_call_id=tool_call_id, turn_id=turn_id, api_request_id=api_request_id, middleware_trace=list(middleware_trace or []), ) - modified_args: Optional[Dict[str, Any]] = None - for result in hook_results: if not isinstance(result, dict): continue @@ -2019,7 +1812,6 @@ def _get_pre_tool_call_directive_details( return _PreToolCallDirective( action=action, message=message, rule_key=rule_key, modified_args=modified_args, ) - return _PreToolCallDirective(modified_args=modified_args) @@ -2063,7 +1855,6 @@ def _resolve_block_from_details( request_tool_approval, reset_current_observability_context, set_current_observability_context, ) - approval_tokens = None with suppress(Exception): approval_tokens = set_current_observability_context( @@ -2111,17 +1902,13 @@ def get_pre_verify_continue_message( "pre_verify", session_id=session_id, platform=platform, model=model, coding=coding, attempt=attempt, final_response=final_response, changed_paths=list(changed_paths or []), ) - for result in hook_results: if not isinstance(result, dict): continue action = str(result.get("action") or result.get("decision") or "").strip().lower() - if action not in ("continue", "block"): - continue message = result.get("message") or result.get("reason") - if isinstance(message, str) and message.strip(): + if action in ("continue", "block") and isinstance(message, str) and message.strip(): return message.strip() - return None @@ -2132,10 +1919,9 @@ def get_plugin_error_classification( num_messages: int = 0, ) -> Optional[Dict[str, Any]]: """Consult ``transform_api_error_classification`` hooks BEFORE the built-in classifier. - Run-all-then-pick-first: every callback runs isolated, the first valid result in registration - order wins, losing valid results warn (conflicts visible, not shadowed). Returns a sanitized dict - (``reason`` -> ``FailoverReason``, hint flags -> bool, ``message`` capped at 500) or ``None``. - Privacy: ``error_message``/``error_body`` may be unredacted.""" + Run-all-then-pick-first: the first valid result in registration order wins, losing valid results + warn (conflicts visible, not shadowed). Returns a sanitized dict (``reason`` -> ``FailoverReason``, + hint flags -> bool, ``message`` capped at 500) or ``None``. Privacy: inputs may be unredacted.""" from agent.error_classifier import FailoverReason hook_results = invoke_hook( "transform_api_error_classification", provider=provider, model=model, @@ -2145,41 +1931,32 @@ def get_plugin_error_classification( num_messages=num_messages, ) - winner: Optional[Dict[str, Any]] = None - skipped_valid = 0 - for result in hook_results: - if not isinstance(result, dict): - continue - reason = result.get("reason") + def _reason(result: Any) -> Any: + reason = result.get("reason") if isinstance(result, dict) else None if isinstance(reason, str): - try: - reason = FailoverReason(reason.strip().lower()) - except ValueError: - continue - if not isinstance(reason, FailoverReason): - continue + with suppress(ValueError): + return FailoverReason(reason.strip().lower()) + return None + return reason if isinstance(reason, FailoverReason) else None - if winner is not None: - skipped_valid += 1 - continue - - out: Dict[str, Any] = {"reason": reason} - for key in ("retryable", "should_compress", "should_rotate_credential", "should_fallback"): - if key in result: - out[key] = bool(result[key]) - message = result.get("message") - if isinstance(message, str) and message.strip(): - out["message"] = message.strip()[:500] - error_context = result.get("error_context") - if isinstance(error_context, dict): - out["error_context"] = error_context - winner = out - - if winner is not None and skipped_valid: + valid = [(result, reason) for result in hook_results if (reason := _reason(result)) is not None] + if not valid: + return None + result, reason = valid[0] + winner: Dict[str, Any] = {"reason": reason} + for key in ("retryable", "should_compress", "should_rotate_credential", "should_fallback"): + if key in result: + winner[key] = bool(result[key]) + message = result.get("message") + if isinstance(message, str) and message.strip(): + winner["message"] = message.strip()[:500] + if isinstance(result.get("error_context"), dict): + winner["error_context"] = result["error_context"] + if len(valid) > 1: logger.warning( "transform_api_error_classification: skipped %d valid " "classification(s) after the first result in registration order " - "won (run-all-then-pick-first)", skipped_valid, + "won (run-all-then-pick-first)", len(valid) - 1, ) return winner @@ -2211,12 +1988,10 @@ def resolve_plugin_command_result(result: Any) -> Any: terminal).""" if not inspect.isawaitable(result): return result - try: asyncio.get_running_loop() except RuntimeError: return asyncio.run(result) - outcome: Dict[str, Any] = {} failure: Dict[str, BaseException] = {} done = threading.Event() @@ -2229,8 +2004,7 @@ def resolve_plugin_command_result(result: Any) -> Any: finally: done.set() - thread = threading.Thread(target=_runner, name="hermes-plugin-command-await", daemon=True) - thread.start() + threading.Thread(target=_runner, name="hermes-plugin-command-await", daemon=True).start() if not done.wait(timeout=_PLUGIN_COMMAND_AWAIT_TIMEOUT_SECS): raise TimeoutError( "Plugin command async handler did not complete within " @@ -2257,27 +2031,23 @@ def get_plugin_toolsets() -> List[tuple]: manager = get_plugin_manager() if not manager._plugin_tool_names: return [] - try: from tools.registry import registry except Exception: return [] - - # Group plugin tool names by their toolset + # Group plugin tool names by their toolset, then map each toolset back to the plugin that + # registered it (first owner wins) for the description. toolset_tools: Dict[str, List[str]] = {} for tool_name in manager._plugin_tool_names: entry = registry.get_entry(tool_name) if entry: toolset_tools.setdefault(entry.toolset, []).append(entry.name) - - # Map toolsets back to the plugin that registered them toolset_plugin: Dict[str, LoadedPlugin] = {} for loaded in manager._plugins.values(): for tool_name in loaded.tools_registered: entry = registry.get_entry(tool_name) if entry and entry.toolset in toolset_tools: toolset_plugin.setdefault(entry.toolset, loaded) - result = [] for ts_key in sorted(toolset_tools): plugin = toolset_plugin.get(ts_key) diff --git a/hermes_cli/plugins_discovery.py b/hermes_cli/plugins_discovery.py index a5c8b54caa..7eba65b277 100644 --- a/hermes_cli/plugins_discovery.py +++ b/hermes_cli/plugins_discovery.py @@ -24,10 +24,7 @@ from hermes_cli.relay_plugin_cutover import LEGACY_RELAY_PLUGIN_KEYS, RELAY_PLUG logger = logging.getLogger("hermes_cli.plugins") - ENTRY_POINTS_GROUP = "hermes_agent.plugins" - - ENTRY_POINT_CAPABILITIES_GROUP = "hermes_agent.plugin_capabilities" @@ -56,30 +53,24 @@ def discover_entrypoint_manifests() -> List["PluginManifest"]: except Exception as exc: logger.debug("Entry-point scan failed: %s", exc) return manifests - for ep in group_eps: try: - capabilities = [] - for capability in VALID_CAPABILITY_IDS: - declaration_name = f"{ep.name}.{capability}" + capabilities = [ + capability for capability in VALID_CAPABILITY_IDS if any( - declaration.name == declaration_name and declaration.value == ep.value + declaration.name == f"{ep.name}.{capability}" and declaration.value == ep.value for declaration in capability_eps - ): - capabilities.append(capability) + ) + ] dist = getattr(ep, "dist", None) metadata = getattr(dist, "metadata", None) - manifest = PluginManifest( + manifests.append(PluginManifest( name=ep.name, version=str(getattr(dist, "version", "") or ""), - description=( - str(metadata.get("Summary", "") or "") if metadata is not None else "" - ), source="entrypoint", path=ep.value, key=ep.name, - capabilities=_parse_declared_capabilities( - capabilities, ep.name - ), - ) - manifest.kind = _classify_entrypoint_value_kind(ep.value) - manifests.append(manifest) + description=str(metadata.get("Summary", "") or "") if metadata is not None else "", + source="entrypoint", path=ep.value, key=ep.name, + kind=_classify_entrypoint_value_kind(ep.value), + capabilities=_parse_declared_capabilities(capabilities, ep.name), + )) except Exception as exc: logger.debug("Entry-point manifest for %r skipped: %s", getattr(ep, "name", "?"), exc) return manifests @@ -100,8 +91,7 @@ def _get_disabled_plugins() -> set: """Read ``plugins.disabled`` — a deny-list that wins over ``plugins.enabled``.""" try: from hermes_cli.config import load_config - config = load_config() - disabled = cfg_get(config, "plugins", "disabled", default=[]) + disabled = cfg_get(load_config(), "plugins", "disabled", default=[]) return set(disabled) if isinstance(disabled, list) else set() except Exception: return set() @@ -115,8 +105,7 @@ def _get_enabled_plugins() -> Optional[set]: """ try: from hermes_cli.config import load_config - plugins_cfg = load_config().get("plugins") - enabled = plugins_cfg.get("enabled") if isinstance(plugins_cfg, dict) else None + enabled = cfg_get(load_config(), "plugins", "enabled") return set(enabled) if isinstance(enabled, list) else None except Exception: return None @@ -134,9 +123,7 @@ def scan_directory( if not path.is_dir(): return manifests for child in sorted(path.iterdir()): - if not child.is_dir(): - continue - if depth == 0 and skip_names and child.name in skip_names: + if not child.is_dir() or (depth == 0 and skip_names and child.name in skip_names): continue manifest_file = child / "plugin.yaml" if not manifest_file.exists(): @@ -167,31 +154,25 @@ def collect_directory_manifests() -> List[PluginManifest]: precedence/containment rules of the real discovery sweep.""" from hermes_cli import plugins as _origin # patched names resolve through the origin manifests: List[PluginManifest] = [] + + def _scan(label: str, directory: Path, source: str, skip_names: Optional[Set[str]] = None) -> None: + found = scan_directory(directory, source, skip_names=skip_names) + logger.debug(" %s: %d manifest(s)", label, len(found)) + manifests.extend(found) + # Excluded bundled top-level categories have their own discovery; platforms scan separately. repo_plugins = _origin.get_bundled_plugins_dir() logger.debug("Scanning bundled plugins: %s", repo_plugins) - bundled = scan_directory( - repo_plugins, "bundled", - skip_names={"memory", "context_engine", "platforms", "model-providers"}, - ) - logger.debug(" bundled (top-level): %d manifest(s)", len(bundled)) - manifests.extend(bundled) - bundled_platforms = scan_directory(repo_plugins / "platforms", "bundled") - logger.debug(" bundled/platforms: %d manifest(s)", len(bundled_platforms)) - manifests.extend(bundled_platforms) - + _scan("bundled (top-level)", repo_plugins, "bundled", + {"memory", "context_engine", "platforms", "model-providers"}) + _scan("bundled/platforms", repo_plugins / "platforms", "bundled") user_dir = get_hermes_home() / "plugins" logger.debug("Scanning user plugins: %s", user_dir) - user_manifests = scan_directory(user_dir, "user") - logger.debug(" user: %d manifest(s)", len(user_manifests)) - manifests.extend(user_manifests) - + _scan("user", user_dir, "user") if _origin._env_enabled("HERMES_ENABLE_PROJECT_PLUGINS"): project_dir = Path.cwd() / ".hermes" / "plugins" logger.debug("Scanning project plugins: %s", project_dir) - project_manifests = scan_directory(project_dir, "project") - logger.debug(" project: %d manifest(s)", len(project_manifests)) - manifests.extend(project_manifests) + _scan("project", project_dir, "project") else: logger.debug("Project plugins disabled (set HERMES_ENABLE_PROJECT_PLUGINS=1 to enable)") return manifests @@ -214,9 +195,9 @@ def gate_manifest( disable, category-owned kinds (exclusive / model-provider), bundled auto-loads (backend now, platform deferred), then ``plugins.enabled`` opt-in (path-derived key or legacy bare name).""" lookup_key = manifest_key(manifest) - name = manifest.name + names = {lookup_key, manifest.name} # Relay lifecycle is core-owned; an old plugin copy would compete for its registries. - if lookup_key in LEGACY_RELAY_PLUGIN_KEYS or name in LEGACY_RELAY_PLUGIN_KEYS: + if names & LEGACY_RELAY_PLUGIN_KEYS: error = ( "removed — Relay lifecycle is owned by Hermes core; configure " f"{RELAY_PLUGINS_CONFIG_ENV} instead" @@ -225,7 +206,7 @@ def gate_manifest( "placeholder", error=error, log=(logging.WARNING, "Refusing to load removed Hermes Relay plugin '%s'; %s", lookup_key, error), ) - if lookup_key in disabled or name in disabled: + if names & disabled: return ManifestGate( "placeholder", error="disabled via config", log=(logging.DEBUG, "Skipping disabled plugin '%s'", lookup_key), @@ -243,14 +224,15 @@ def gate_manifest( "placeholder", enabled=True, log=(logging.DEBUG, "Skipping '%s' (model-provider, handled by providers/ discovery)", lookup_key), ) - # Bundled backends auto-load; selection among them is ``.provider`` config. - if manifest.source == "bundled" and manifest.kind == "backend": - return ManifestGate("load_now") - # Bundled platforms register LAZILY: eagerly importing ~20 heavy SDKs added seconds to every - # `hermes` invocation. A deferred loader keeps every platform available on first use. - if manifest.source == "bundled" and manifest.kind == "platform": - return ManifestGate("defer") - if enabled is None or not (lookup_key in enabled or name in enabled): + if manifest.source == "bundled": + # Bundled backends auto-load; selection among them is ``.provider`` config. + if manifest.kind == "backend": + return ManifestGate("load_now") + # Bundled platforms register LAZILY: eagerly importing ~20 heavy SDKs added seconds to every + # `hermes` invocation. A deferred loader keeps every platform available on first use. + if manifest.kind == "platform": + return ManifestGate("defer") + if enabled is None or not names & enabled: return ManifestGate( "placeholder", error=f"not enabled in config (run `hermes plugins enable {lookup_key}` to activate)", diff --git a/hermes_cli/plugins_ledger.py b/hermes_cli/plugins_ledger.py index 33139828fb..9489102ee1 100644 --- a/hermes_cli/plugins_ledger.py +++ b/hermes_cli/plugins_ledger.py @@ -59,18 +59,16 @@ class PluginLedgerMixin: self, manifest: PluginManifest, kind: str, key: str, release: Callable[[], None], *, persistent: bool = False, ) -> PluginRegistration: - """Record one registration under its canonical plugin key. - - ``persistent`` ones (process-global host infrastructure) stay in the ownership ledger for - attribution but NOT in ``_registration_order``, so a routine unload cannot dispose them; - the handle still releases on explicit ``dispose()``. - """ - plugin_key = manifest_key(manifest) + """Record one registration under its canonical plugin key. ``persistent`` ones + (process-global host infrastructure) stay in the ownership ledger for attribution but NOT in + ``_registration_order``, so a routine unload cannot dispose them; the handle still releases + on explicit ``dispose()``.""" registration = PluginRegistration( - kind=kind, key=key, release=release, plugin_key=plugin_key, persistent=persistent, + kind=kind, key=key, release=release, plugin_key=manifest_key(manifest), + persistent=persistent, ) registration._on_dispose = lambda disposed: self._forget_registrations([disposed]) - self._ownership_ledger.setdefault(plugin_key, []).append(registration) + self._ownership_ledger.setdefault(registration.plugin_key, []).append(registration) if not persistent: self._registration_order.append(registration) return registration @@ -79,11 +77,9 @@ class PluginLedgerMixin: self, manifest: PluginManifest, kind: str, name: str, registry: Any, current: Any, previous: Any, *, finalize: Optional[Callable[[], None]] = None, ) -> PluginRegistration: - """Lease one ``(kind, scope, name)`` slot of a scope-keyed process-global registry. - - Unload calls ``registry.restore_registration(name, current, replacement, scope=...)`` — - identity-conditional, so a later generation is never removed by an earlier owner. - """ + """Lease one ``(kind, scope, name)`` slot of a scope-keyed process-global registry. Unload + calls ``registry.restore_registration(name, current, replacement, scope=...)`` — + identity-conditional, so a later generation is never removed by an earlier owner.""" scope = self.scope_key lease = replacement_coordinator.acquire( (kind, scope, name), current=current, previous=previous, @@ -93,22 +89,22 @@ class PluginLedgerMixin: ) return self._track_registration(manifest, kind, name, lease.dispose) + def _active_persistent(self) -> List[PluginRegistration]: + """Live persistent registrations across every plugin in the ownership ledger.""" + return [ + registration for owned in self._ownership_ledger.values() + for registration in owned if registration.persistent and registration.active + ] + def _evict_stale_persistent_registrations(self) -> None: """After re-discovery, dispose parked persistent handles whose plugin did not re-register the same ``(kind, key)``. Re-registered ones are dropped WITHOUT disposing — the same object re-registered would pass the identity check and evict the live entry.""" if not self._persistent_carryover: return - parked = self._persistent_carryover - self._persistent_carryover = [] - current = { - (registration.kind, registration.key) for owned in self._ownership_ledger.values() - for registration in owned if registration.persistent and registration.active - } - stale = [ - registration for registration in parked if registration.active - and (registration.kind, registration.key) not in current - ] + parked, self._persistent_carryover = self._persistent_carryover, [] + current = {(r.kind, r.key) for r in self._active_persistent()} + stale = [r for r in parked if r.active and (r.kind, r.key) not in current] for registration in stale: logger.info( "Evicting persistent registration %s/%s: plugin '%s' no " @@ -158,8 +154,7 @@ class PluginLedgerMixin: def _remove_name_if_unowned(self, kind: str, names: Set[str], name: str) -> None: """Drop *name* from the manager-local name set once no active ledger entry owns it.""" if not any( - registration.active and registration.kind == kind and registration.key == name - for registration in self._registration_order + r.active and r.kind == kind and r.key == name for r in self._registration_order ): names.discard(name) @@ -172,15 +167,10 @@ class PluginLedgerMixin: def _forget_registrations(self, registrations: List[PluginRegistration]) -> None: if not registrations: return - registration_ids = {id(registration) for registration in registrations} - self._registration_order = [ - registration for registration in self._registration_order - if id(registration) not in registration_ids - ] + ids = {id(r) for r in registrations} + self._registration_order = [r for r in self._registration_order if id(r) not in ids] for plugin_key, owned in list(self._ownership_ledger.items()): - remaining = [ - registration for registration in owned if id(registration) not in registration_ids - ] + remaining = [r for r in owned if id(r) not in ids] if remaining: self._ownership_ledger[plugin_key] = remaining else: @@ -204,7 +194,7 @@ class PluginLedgerMixin: if isinstance(plugin, LoadedPlugin): return manifest_key(plugin.manifest) if isinstance(plugin, PluginManifest): - return plugin.key or plugin.name + return manifest_key(plugin) return str(plugin) def unload(self, plugin: Union[str, PluginManifest, LoadedPlugin, None] = None) -> bool: @@ -213,41 +203,32 @@ class PluginLedgerMixin: return self._unload_scoped(plugin) def _unload_scoped(self, plugin: Union[str, PluginManifest, LoadedPlugin, None] = None) -> bool: - """Unload one plugin (or all when ``plugin=None``, as force rediscovery does). - - Every ledger registration — including on_unload callbacks and supervised tasks — is disposed - in reverse acquisition order with identity-conditional inverses. Returns ``True`` when - anything was found. - """ + """Unload one plugin (or all when ``plugin=None``, as force rediscovery does). Every ledger + registration — including on_unload callbacks and supervised tasks — is disposed in reverse + acquisition order with identity-conditional inverses. Returns ``True`` when anything was + found.""" unload_all = plugin is None if unload_all: target_keys = set(self._ownership_ledger) | set(self._plugins) registrations = list(self._registration_order) else: target_keys = self._unload_target_keys(self._resolve_plugin_key(plugin)) - registrations = [ - registration for registration in self._registration_order - if registration.plugin_key in target_keys - ] + registrations = [r for r in self._registration_order if r.plugin_key in target_keys] # Persistent registrations are absent from _registration_order (unload-all keeps them), # but a *targeted* unload is the disable/uninstall path: a disabled auth plugin's # provider must NOT stay live process-wide. registrations.extend( - registration for key in target_keys - for registration in self._ownership_ledger.get(key, []) - if registration.persistent and registration.active + r for key in target_keys for r in self._ownership_ledger.get(key, []) + if r.persistent and r.active ) - found = bool(target_keys or registrations) self._dispose_registrations(registrations) self._forget_registrations(registrations) - if unload_all: self._reset_after_unload_all(registrations) else: for key in target_keys: self._plugins.pop(key, None) - return found def _unload_target_keys(self, requested: str) -> Set[str]: @@ -266,12 +247,8 @@ class PluginLedgerMixin: platform_registry.unregister(platform_name) # Ledger-owned tool names are excluded: their handles already restored the previous entry, # and blanket deregistration would remove what the ledger just restored. - ledger_tool_names = { - registration.key for registration in registrations if registration.kind == "tool" - } - preledger_tools = tuple( - name for name in self._plugin_tool_names if name not in ledger_tool_names - ) + ledger_tool_names = {r.key for r in registrations if r.kind == "tool"} + preledger_tools = tuple(n for n in self._plugin_tool_names if n not in ledger_tool_names) if preledger_tools: try: from tools.registry import registry as tool_registry @@ -285,11 +262,9 @@ class PluginLedgerMixin: logger.debug("unload: tool deregister %s failed: %s", tool_name, exc) # Persistent registrations survive unload-all but must not be orphaned by the ledger clear: # carry them over so force re-discovery can evict the ones whose plugin does not come back. - carryover_ids = {id(registration) for registration in self._persistent_carryover} + carryover_ids = {id(r) for r in self._persistent_carryover} self._persistent_carryover.extend( - registration for owned in self._ownership_ledger.values() for registration in owned - if registration.persistent and registration.active - and id(registration) not in carryover_ids + r for r in self._active_persistent() if id(r) not in carryover_ids ) for container in ( self._ownership_ledger, self._plugins, self._hooks, self._middleware, diff --git a/hermes_cli/plugins_loader.py b/hermes_cli/plugins_loader.py index 6b26502f96..1abf7c8887 100644 --- a/hermes_cli/plugins_loader.py +++ b/hermes_cli/plugins_loader.py @@ -25,20 +25,16 @@ from hermes_constants import get_hermes_home, reset_hermes_home_override, set_he from registration_lifecycle import replacement_coordinator from hermes_cli.plugins_discovery import ENTRY_POINTS_GROUP, _select_entry_point_group from hermes_cli.plugins_manifest import PluginManifest, manifest_key, validate_config_schema +from hermes_cli.plugins_state import _plugin_settings_entry if TYPE_CHECKING: # pragma: no cover from hermes_cli.plugins import LoadedPlugin logger = logging.getLogger("hermes_cli.plugins") - _NS_PARENT = "hermes_plugins" - - _MODULE_NAMESPACE_LOCK = threading.RLock() - - -_BARE_MODULE_SCOPE: Dict[str, str] = {} +_BARE_MODULE_SCOPE: Dict[str, str] = {} # bare module name -> owning scope_key def _evict_modules(module_name: str) -> None: @@ -69,15 +65,14 @@ def _plugin_home_scope(home: Path): class PluginLoaderMixin: - def _platform_name_from_manifest(self, manifest: PluginManifest) -> str: + @staticmethod + def _platform_name_from_manifest(manifest: PluginManifest) -> str: """Derive the platform name without importing the adapter: strip a trailing ``-platform`` from the manifest name, else the directory basename (the bundled convention).""" name = manifest.name or "" if name.endswith("-platform"): return name[: -len("-platform")] - if manifest.path: - return Path(manifest.path).name - return name + return Path(manifest.path).name if manifest.path else name @_serialized_replacement def _register_deferred_platform(self, manifest: PluginManifest) -> None: @@ -87,14 +82,10 @@ class PluginLoaderMixin: from hermes_cli.plugins import LoadedPlugin lookup_key = manifest_key(manifest) platform_name = self._platform_name_from_manifest(manifest) - - loaded = LoadedPlugin(manifest=manifest, enabled=True) - loaded.deferred = True + loaded = LoadedPlugin(manifest=manifest, enabled=True, deferred=True) self._plugins[lookup_key] = loaded - try: from gateway.platform_registry import platform_registry - scope = self.scope_key def _loader(_manifest: PluginManifest = manifest) -> None: @@ -126,81 +117,75 @@ class PluginLoaderMixin: ) self._load_plugin(manifest) return - self._register_deferred_platform_tools(manifest, loaded) def _register_deferred_platform_tools( self, manifest: PluginManifest, loaded: LoadedPlugin ) -> None: - """Register a deferred platform's *client* tools without its adapter. - - Deferring the plugin would otherwise defer its outbound tools too, so in CLI/TUI processes - (which never materialize platforms) they would be missing from ``hermes tools`` and dropped - from ``platform_toolsets``. Opt-in is explicit via ``provides_tools``; tools live in a - ``tools`` submodule so ``__init__`` stays import-light. - """ + """Register a deferred platform's *client* tools without its adapter. Deferring the plugin + would otherwise defer its outbound tools too, so CLI/TUI processes (which never materialize + platforms) would miss them in ``hermes tools`` / ``platform_toolsets``. Opt-in is explicit via + ``provides_tools``; tools live in a ``tools`` submodule so ``__init__`` stays import-light.""" from hermes_cli.plugins import PluginContext, _PLUGINS_DEBUG if not manifest.provides_tools: return - lookup_key = manifest_key(manifest) + declared = list(manifest.provides_tools) plugin_dir = Path(manifest.path) if manifest.path else None if plugin_dir is None or not (plugin_dir / "tools.py").is_file(): # Declared but undeliverable — staying quiet reproduces the very symptom this fixes. logger.warning( "Plugin '%s' declares provides_tools %s but has no tools.py; " - "those tools will not be available in CLI/TUI sessions.", lookup_key, - list(manifest.provides_tools), + "those tools will not be available in CLI/TUI sessions.", lookup_key, declared, ) return - before = set(self._plugin_tool_names) # lets the failure path credit partial registrations + + def _credit() -> List[str]: + """Attribute every tool registered since ``before`` to this plugin.""" + registered = [t for t in self._plugin_tool_names if t not in before] + if registered: + loaded.tools_registered = registered + self._predeclared_tools[lookup_key] = registered + return registered + try: module = self._load_directory_module(manifest) # Record the module even if nothing registers: the package body has run, so # materializing the adapter later must reuse it rather than execute it twice. loaded.module = module self._predeclared_modules[lookup_key] = module - - tools_module = importlib.import_module(f"{module.__name__}.tools") - register_tools = getattr(tools_module, "register_tools", None) + register_tools = getattr(importlib.import_module(f"{module.__name__}.tools"), + "register_tools", None) if register_tools is None: logger.warning( "Plugin '%s' declares provides_tools %s but its tools.py " "has no register_tools(ctx); those tools will not be " - "available in CLI/TUI sessions.", lookup_key, list(manifest.provides_tools), + "available in CLI/TUI sessions.", lookup_key, declared, ) return - register_tools(PluginContext(manifest, self)) - registered = [t for t in self._plugin_tool_names if t not in before] - - loaded.tools_registered = registered - self._predeclared_tools[lookup_key] = registered + registered = _credit() logger.debug( "Deferred platform '%s': pre-registered %d client tool(s) %s", lookup_key, len(registered), registered, ) except Exception as exc: # Tools registered before the raise are live: credit them or `hermes plugins list` - # under-reports (and _load_plugin's later diff would miss them too). - partial = [t for t in self._plugin_tool_names if t not in before] - if partial: - loaded.tools_registered = partial - self._predeclared_tools[lookup_key] = partial - - # Never break discovery (the platform stays deferred), but a broken tools.py IS the - # symptom, so warn — and say where it failed, which is what the operator needs first. - declared = len(manifest.provides_tools) + # under-reports (and _load_plugin's later diff would miss them too). Never break + # discovery (the platform stays deferred), but a broken tools.py IS the symptom, so warn + # — and say where it failed, which is what the operator needs first. + partial = _credit() + total = len(declared) if not partial: - scope = f"before registering any of its {declared} declared tool(s)" - elif len(partial) >= declared: - scope = f"after registering all {declared} declared tool(s)" + scope = f"before registering any of its {total} declared tool(s)" + elif len(partial) >= total: + scope = f"after registering all {total} declared tool(s)" else: - scope = f"after registering {len(partial)} of {declared} declared tool(s)" + scope = f"after registering {len(partial)} of {total} declared tool(s)" logger.warning( "Plugin '%s': client-tool pre-registration failed %s (%s).%s", lookup_key, scope, - exc, "" if len(partial) >= declared else + exc, "" if len(partial) >= total else " The remainder will be missing from CLI/TUI sessions.", exc_info=_PLUGINS_DEBUG, ) @@ -240,14 +225,10 @@ class PluginLoaderMixin: settings: Mapping[str, Any] = {} try: from hermes_cli.config import load_config - - cfg = load_config() or {} - entries = (cfg.get("plugins") or {}).get("entries") or {} - entry = entries.get(plugin_id) if isinstance(entries, Mapping) else None - raw = entry.get("settings") if isinstance(entry, Mapping) else None + entry = _plugin_settings_entry(load_config() or {}, plugin_id) or {} + raw = entry.get("settings") if not isinstance(raw, Mapping): - # Migration fallback mirroring ctx.get_config. - raw = entry.get("config") if isinstance(entry, Mapping) else None + raw = entry.get("config") # migration fallback mirroring ctx.get_config settings = raw if isinstance(raw, Mapping) else {} except Exception: settings = {} @@ -263,29 +244,25 @@ class PluginLoaderMixin: """Load one plugin with the manager's home bound as current.""" from hermes_cli.plugins import LoadedPlugin, PluginContext, _PLUGINS_DEBUG loaded = LoadedPlugin(manifest=manifest) + plugin_key = manifest_key(manifest) logger.debug( "Loading plugin '%s' (source=%s, kind=%s, path=%s)", - manifest_key(manifest), manifest.source, manifest.kind, manifest.path, + plugin_key, manifest.source, manifest.kind, manifest.path, ) - if manifest.portable: self._load_portable_plugin(manifest, loaded) return - registration_start = len(self._registration_order) - plugin_key = manifest_key(manifest) - _module_name = self._policy_module_name(manifest) - self._track_tool_override_policy(manifest, _module_name) + module_name = self._policy_module_name(manifest) + self._track_tool_override_policy(manifest, module_name) try: # Reuse a deferred platform's already-imported package so its body doesn't run twice. - preloaded = self._predeclared_modules.pop(plugin_key, None) - if preloaded is not None: - module = preloaded - elif manifest.source in {"user", "project", "bundled"}: - module = self._load_directory_module(manifest, module_name=_module_name) - else: - module = self._load_entrypoint_module(manifest) - + module = self._predeclared_modules.pop(plugin_key, None) + if module is None: + if manifest.source in {"user", "project", "bundled"}: + module = self._load_directory_module(manifest, module_name=module_name) + else: + module = self._load_entrypoint_module(manifest) loaded.module = module register_fn = getattr(module, "register", None) if register_fn is None: @@ -295,12 +272,8 @@ class PluginLoaderMixin: register_fn(PluginContext(manifest, self)) self._attribute_registrations(loaded, plugin_key, registration_start) loaded.enabled = True - except Exception as exc: - owned = [ - registration for registration in self._registration_order - if registration.plugin_key == plugin_key - ] + owned = [r for r in self._registration_order if r.plugin_key == plugin_key] self._dispose_registrations(owned) self._forget_registrations(owned) loaded.error = str(exc) @@ -314,26 +287,23 @@ class PluginLoaderMixin: # so discovery-time pre-registrations are gone too. if not loaded.enabled: self._predeclared_tools.pop(plugin_key, None) - self._plugins[manifest_key(manifest)] = loaded + self._plugins[plugin_key] = loaded def _track_tool_override_policy(self, manifest: PluginManifest, module_name: str) -> None: """Install the plugin's tool-override policy in tools.registry as a ledger-owned lease.""" from hermes_cli.plugins import PluginContext - from tools.registry import registry as _registry + scope = self.scope_key with replacement_coordinator.transaction(): - previous_policy = _registry.snapshot_plugin_override_policy( - module_name, scope=self.scope_key - ) + previous_policy = _registry.snapshot_plugin_override_policy(module_name, scope=scope) current_policy = _registry.register_plugin_override_policy( - module_name, PluginContext(manifest, self)._tool_override_allowed(""), - scope=self.scope_key, + module_name, PluginContext(manifest, self)._tool_override_allowed(""), scope=scope, ) policy_lease = replacement_coordinator.acquire( - ("tool_override_policy", self.scope_key, module_name), current=current_policy, + ("tool_override_policy", scope, module_name), current=current_policy, previous=previous_policy, restore=lambda replacement: _registry.restore_plugin_override_policy( - module_name, current_policy, replacement, scope=self.scope_key, + module_name, current_policy, replacement, scope=scope, ), ) self._track_registration( @@ -345,8 +315,8 @@ class PluginLoaderMixin: ) -> None: """Fill ``loaded.*_registered`` from the ledger slice this plugin's register() produced.""" registrations = [ - registration for registration in self._registration_order[registration_start:] - if registration.plugin_key == plugin_key and registration.active + r for r in self._registration_order[registration_start:] + if r.plugin_key == plugin_key and r.active ] def _keys(kind: str) -> List[str]: @@ -354,12 +324,10 @@ class PluginLoaderMixin: # Discovery-time tools predate registration_start; credit them back or `hermes plugins # list` under-reports once the deferred adapter materializes. - _predeclared = [ + predeclared = [ t for t in self._predeclared_tools.pop(plugin_key, []) if t in self._plugin_tool_names ] - loaded.tools_registered = _predeclared + [ - key for key in _keys("tool") if key not in _predeclared - ] + loaded.tools_registered = predeclared + [k for k in _keys("tool") if k not in predeclared] loaded.hooks_registered = _keys("hook") loaded.middleware_registered = _keys("middleware") loaded.commands_registered = _keys("command") @@ -376,7 +344,6 @@ class PluginLoaderMixin: lookup_key = manifest_key(manifest) try: from hermes_cli.agent_plugins import load_agent_plugin - package = load_agent_plugin( Path(manifest.path), get_hermes_home() / "plugin-data" / manifest.skill_namespace, ) @@ -409,15 +376,12 @@ class PluginLoaderMixin: self._plugins[lookup_key] = loaded def _directory_module_name(self, manifest: PluginManifest) -> str: - """Return a profile-safe import namespace for a directory plugin.""" - key = manifest_key(manifest) - slug = key.replace("/", "__").replace("-", "_") + """Profile-safe import namespace for a directory plugin: the bare ``hermes_plugins.`` + for the first scope that claims it, a ``__home_`` suffix for any other scope.""" + slug = manifest_key(manifest).replace("/", "__").replace("-", "_") bare_name = f"{_NS_PARENT}.{slug}" with _MODULE_NAMESPACE_LOCK: - owner = _BARE_MODULE_SCOPE.get(bare_name) - if owner is None: - _BARE_MODULE_SCOPE[bare_name] = self.scope_key - return bare_name + owner = _BARE_MODULE_SCOPE.setdefault(bare_name, self.scope_key) if owner == self.scope_key: return bare_name digest = hashlib.sha256(self.scope_key.encode("utf-8")).hexdigest()[:12] @@ -440,27 +404,22 @@ class PluginLoaderMixin: init_file = plugin_dir / "__init__.py" if not init_file.exists(): raise FileNotFoundError(f"No __init__.py in {plugin_dir}") - if _NS_PARENT not in sys.modules: ns_pkg = types.ModuleType(_NS_PARENT) ns_pkg.__path__ = [] # type: ignore[attr-defined] ns_pkg.__package__ = _NS_PARENT sys.modules[_NS_PARENT] = ns_pkg - module_name = module_name or self._directory_module_name(manifest) - # Evict stale entries for this slug (same slug cached from another Hermes home, or an # earlier force reload). Replacing only sys.modules[module_name] is not enough: the plugin's # relative imports are cached as "module_name.sub" and resolve from sys.modules first, so a # stale submodule would keep serving the previous load's code/state. _evict_modules(module_name) - spec = importlib.util.spec_from_file_location( module_name, init_file, submodule_search_locations=[str(plugin_dir)], ) if spec is None or spec.loader is None: raise ImportError(f"Cannot create module spec for {init_file}") - module = importlib.util.module_from_spec(spec) module.__package__ = module_name module.__path__ = [str(plugin_dir)] # type: ignore[attr-defined] @@ -479,7 +438,6 @@ class PluginLoaderMixin: for ep in _select_entry_point_group(importlib.metadata.entry_points(), ENTRY_POINTS_GROUP): if ep.name == manifest.name: return ep.load() - raise ImportError( f"Entry point '{manifest.name}' not found in group '{ENTRY_POINTS_GROUP}'" ) diff --git a/hermes_cli/plugins_manifest.py b/hermes_cli/plugins_manifest.py index d2b0e5e2ce..545076a6d8 100644 --- a/hermes_cli/plugins_manifest.py +++ b/hermes_cli/plugins_manifest.py @@ -23,22 +23,38 @@ except ImportError: # pragma: no cover – yaml is optional at import time logger = logging.getLogger("hermes_cli.plugins") +_VALID_PLUGIN_KINDS: Set[str] = {"standalone", "backend", "exclusive", "platform", "model-provider"} + +# Unknown plugin.yaml fields are forward-compat surface: warn (debug for v1 files, warning for v2+) +# and continue loading. ``capabilities``/``emits``/``listens``/``hermes``/``depends`` are reserved. +_KNOWN_MANIFEST_FIELDS: Set[str] = { + "name", "version", "description", "author", "requires_env", "provides_tools", "provides_hooks", + "kind", "hooks", "label", "optional_env", "platforms", "external_dependencies", + "pip_dependencies", "provides_browser_providers", "provides_web_providers", + "manifest_version", "api_version", "requires_plugins", "python_dependencies", "config_schema", + "license", "homepage", "tags", "capabilities", "emits", "listens", "hermes", "depends", +} + +# Highest manifest schema version this Hermes understands. +SUPPORTED_MANIFEST_VERSION = 2 + +_CONFIG_SCHEMA_TYPES: Dict[str, tuple] = { + "str": (str,), "string": (str,), "int": (int,), "integer": (int,), "float": (int, float), + "number": (int, float), "bool": (bool,), "boolean": (bool,), "list": (list,), "array": (list,), + "dict": (dict,), "object": (dict,), +} + def _plugins_debug() -> bool: from hermes_cli import plugins as _origin return _origin._PLUGINS_DEBUG -_VALID_PLUGIN_KINDS: Set[str] = {"standalone", "backend", "exclusive", "platform", "model-provider"} - - def _portable_skill_namespace(key: str) -> str: """Return a readable, collision-resistant namespace for a portable plugin.""" - slug = "".join( ch if ch.isascii() and (ch.isalnum() or ch in "_-") else "-" for ch in key.lower() - ) - slug = slug.strip("-_") or "plugin" + ).strip("-_") or "plugin" digest = hashlib.sha256(key.encode("utf-8")).hexdigest()[:8] return f"agent-plugin-{slug}-{digest}" @@ -46,39 +62,10 @@ def _portable_skill_namespace(key: str) -> str: def _display_author(value: object) -> str: """Normalize a manifest author value for the string PluginManifest field.""" if isinstance(value, Mapping): - return ", ".join( - str(value[field]) for field in ("name", "email", "url") if value.get(field) - ) + return ", ".join(str(value[f]) for f in ("name", "email", "url") if value.get(f)) return "" if value is None else str(value) -# Manifest v2 parsing. Unknown plugin.yaml fields are forward-compat surface: warn (debug for v1 -# files, warning for v2+) and continue loading. -_KNOWN_MANIFEST_FIELDS: Set[str] = { - # v1 - "name", "version", "description", "author", "requires_env", - "provides_tools", "provides_hooks", "kind", "hooks", "label", - "optional_env", "platforms", "external_dependencies", "pip_dependencies", - "provides_browser_providers", "provides_web_providers", - # v2 - "manifest_version", "api_version", "requires_plugins", - "python_dependencies", "config_schema", "license", "homepage", "tags", - # owned by sibling sub-issues but reserved so their manifests don't warn - "capabilities", "emits", "listens", "hermes", "depends", -} - - -# Highest manifest schema version this Hermes understands. -SUPPORTED_MANIFEST_VERSION = 2 - - -_CONFIG_SCHEMA_TYPES: Dict[str, tuple] = { - "str": (str,), "string": (str,), "int": (int,), "integer": (int,), "float": (int, float), - "number": (int, float), "bool": (bool,), "boolean": (bool,), "list": (list,), "array": (list,), - "dict": (dict,), "object": (dict,), -} - - def _manifest_field_of_type(data: Mapping, key: str, field_name: str, typ, what: str): """Return ``data[field_name]`` when absent or of ``typ``; warn and return None otherwise.""" raw = data.get(field_name) @@ -91,7 +78,6 @@ def _manifest_field_of_type(data: Mapping, key: str, field_name: str, typ, what: def _parse_manifest_v2_fields(data: Mapping, key: str) -> Dict[str, Any]: """Validate/normalize manifest v2 fields into PluginManifest kwargs (warnings, never failures).""" out: Dict[str, Any] = {} - # manifest_version — absent means v1 (supported forever). raw_mv = data.get("manifest_version", 1) try: @@ -108,7 +94,6 @@ def _parse_manifest_v2_fields(data: Mapping, key: str) -> Dict[str, Any]: key, mv, SUPPORTED_MANIFEST_VERSION, ) out["manifest_version"] = mv - # api_version — plugin API generation (independent of manifest_version). raw_api = data.get("api_version") out["api_version"] = None @@ -117,7 +102,6 @@ def _parse_manifest_v2_fields(data: Mapping, key: str) -> Dict[str, Any]: out["api_version"] = int(raw_api) except (TypeError, ValueError): logger.warning("Plugin %s: api_version %r is not an integer; ignoring", key, raw_api) - # requires_plugins — list of {id, version_range?} (str shorthand ok). deps: List[Dict[str, Any]] = [] for item in _manifest_field_of_type(data, key, "requires_plugins", list, "a list") or []: @@ -132,7 +116,6 @@ def _parse_manifest_v2_fields(data: Mapping, key: str) -> Dict[str, Any]: "string or a {id, version_range} mapping; skipping", key, item, ) out["requires_plugins"] = deps - # python_dependencies — validated and surfaced ONLY; never auto-installed. pydeps: List[str] = [] raw_pydeps = _manifest_field_of_type( @@ -147,7 +130,6 @@ def _parse_manifest_v2_fields(data: Mapping, key: str) -> Dict[str, Any]: "requirement string; skipping", key, item, ) out["python_dependencies"] = pydeps - # config_schema — mapping of key -> {type?, default?, description?, required?}. schema: Dict[str, Any] = {} raw_schema = _manifest_field_of_type(data, key, "config_schema", Mapping, "a mapping") @@ -167,13 +149,10 @@ def _parse_manifest_v2_fields(data: Mapping, key: str) -> Dict[str, Any]: ) schema[str(skey)] = dict(spec) out["config_schema"] = schema - - # Standard metadata. out["license"] = str(data.get("license") or "") out["homepage"] = str(data.get("homepage") or "") raw_tags = _manifest_field_of_type(data, key, "tags", list, "a list") out["tags"] = [str(t) for t in (raw_tags or [])] - # Forward compat: unknown fields warn (never fail); v1 manifests only at debug. unknown = sorted(set(data.keys()) - _KNOWN_MANIFEST_FIELDS) if unknown: @@ -182,7 +161,6 @@ def _parse_manifest_v2_fields(data: Mapping, key: str) -> Dict[str, Any]: "Plugin %s: unknown manifest field(s) ignored: %s " "(newer manifest schema or typo; plugin still loads)", key, ", ".join(unknown), ) - return out @@ -194,8 +172,7 @@ def validate_config_schema(plugin_id: str, schema: Mapping, settings: Mapping) - for skey, spec in schema.items(): if not isinstance(spec, Mapping): continue - present = skey in settings - if not present: + if skey not in settings: if spec.get("required") and "default" not in spec: warnings.append( f"plugins.entries.{plugin_id}.settings.{skey} is required " @@ -204,17 +181,15 @@ def validate_config_schema(plugin_id: str, schema: Mapping, settings: Mapping) - continue stype = spec.get("type") expected = _CONFIG_SCHEMA_TYPES.get(str(stype).lower()) if stype else None - if expected is not None: - value = settings[skey] - # bool is an int subclass — don't let True satisfy int/float. - ok = isinstance(value, expected) and not ( - isinstance(value, bool) and bool not in expected + if expected is None: + continue + value = settings[skey] + # bool is an int subclass — don't let True satisfy int/float. + if not isinstance(value, expected) or (isinstance(value, bool) and bool not in expected): + warnings.append( + f"plugins.entries.{plugin_id}.settings.{skey} should be " + f"{stype} (got {type(value).__name__})" ) - if not ok: - warnings.append( - f"plugins.entries.{plugin_id}.settings.{skey} should be " - f"{stype} (got {type(value).__name__})" - ) return warnings @@ -228,22 +203,15 @@ def resolve_plugin_load_order(manifests: Mapping[str, "PluginManifest"]) -> List keys = sorted(manifests.keys()) by_name: Dict[str, str] = {} for k in keys: - name = manifests[k].name - if name and name not in by_name: - by_name[name] = k - - def _resolve_dep(dep_id: str) -> Optional[str]: - if dep_id in manifests: - return dep_id - return by_name.get(dep_id) - + if manifests[k].name: + by_name.setdefault(manifests[k].name, k) edges: Dict[str, Set[str]] = {k: set() for k in keys} for k in keys: for dep in manifests[k].requires_plugins: dep_id = dep.get("id") if isinstance(dep, Mapping) else None if not dep_id: continue - resolved = _resolve_dep(dep_id) + resolved = dep_id if dep_id in manifests else by_name.get(dep_id) if resolved is None: logger.warning( "Plugin %s requires plugin '%s' which is not enabled/" @@ -251,12 +219,10 @@ def resolve_plugin_load_order(manifests: Mapping[str, "PluginManifest"]) -> List "via ctx.has_plugin). Run `hermes plugins enable %s` if it is installed.", k, dep_id, dep_id, ) - continue - if resolved == k: + elif resolved == k: logger.warning("Plugin %s declares a dependency on itself; ignoring", k) - continue - edges[k].add(resolved) - + else: + edges[k].add(resolved) sorter = graphlib.TopologicalSorter(edges) try: sorter.prepare() @@ -267,7 +233,6 @@ def resolve_plugin_load_order(manifests: Mapping[str, "PluginManifest"]) -> List "alphabetical load order for all plugins", " -> ".join(str(c) for c in cycle), ) return keys - ordered: List[str] = [] while sorter.is_active(): ready = sorted(sorter.get_ready()) @@ -277,11 +242,9 @@ def resolve_plugin_load_order(manifests: Mapping[str, "PluginManifest"]) -> List def _detect_kind_from_source(source_text: str) -> Optional[str]: - """Return the kind implied by source markers (mirrors plugins/memory ``_is_memory_provider_dir``). - - Memory-provider markers -> ``exclusive``; ``register_provider`` + ``ProviderProfile`` -> - ``model-provider``; else ``None``. Keeps both kinds out of the general manager's eager import. - """ + """Kind implied by source markers (mirrors plugins/memory ``_is_memory_provider_dir``): + memory-provider markers -> ``exclusive``; ``register_provider`` + ``ProviderProfile`` -> + ``model-provider``; else ``None``. Keeps both kinds out of the general manager's eager import.""" if "register_memory_provider" in source_text or "MemoryProvider" in source_text: return "exclusive" if "register_provider" in source_text and "ProviderProfile" in source_text: @@ -293,14 +256,11 @@ def _read_source_from_origin(origin: Optional[str], limit: int = 8192) -> str: """First ``limit`` chars of a module's source (``.pyc`` mapped back to ``.py``); "" on failure.""" if not origin: return "" - if origin.endswith((".pyc", ".pyo")): - try: - origin = importlib.util.source_from_cache(origin) - except Exception: - return "" - if not origin.endswith(".py"): - return "" try: + if origin.endswith((".pyc", ".pyo")): + origin = importlib.util.source_from_cache(origin) + if not origin.endswith(".py"): + return "" return Path(origin).read_text(encoding="utf-8", errors="replace")[:limit] except Exception: return "" @@ -322,7 +282,6 @@ def resolve_module_origin(module_name: str) -> Optional[str]: return None if len(parts) == 1: return spec.origin - search_paths = spec.submodule_search_locations if not search_paths: return None diff --git a/hermes_cli/plugins_state.py b/hermes_cli/plugins_state.py index f8d08adc73..1699f7009a 100644 --- a/hermes_cli/plugins_state.py +++ b/hermes_cli/plugins_state.py @@ -13,22 +13,11 @@ from typing import Any, Dict, Mapping from hermes_constants import get_hermes_home from hermes_cli.plugins_manifest import _portable_skill_namespace - _PLUGIN_SETTING_SEGMENT_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$") - - _PLUGIN_SETTING_RESERVED_ROOTS = frozenset({"model", "plugins", "security", "settings"}) - - _PLUGIN_STATE_KEY_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.:-]{0,127}$") - - _PLUGIN_STATE_QUOTA_BYTES = 10 * 1024 * 1024 - - _PLUGIN_STATE_LOCKS: Dict[str, threading.RLock] = {} - - _PLUGIN_STATE_LOCKS_GUARD = threading.Lock() @@ -40,9 +29,8 @@ def _plugin_relative_segments(key: str) -> tuple[str, ...]: segments = tuple(key.split(".")) if ( not key or "/" in key or "\\" in key - or any( - not _PLUGIN_SETTING_SEGMENT_RE.fullmatch(segment) for segment in segments - ) or segments[0].lower() in _PLUGIN_SETTING_RESERVED_ROOTS + or not all(_PLUGIN_SETTING_SEGMENT_RE.fullmatch(segment) for segment in segments) + or segments[0].lower() in _PLUGIN_SETTING_RESERVED_ROOTS ): raise ValueError( "Expected a plugin-relative config key such as 'endpoint' or " @@ -52,6 +40,7 @@ def _plugin_relative_segments(key: str) -> tuple[str, ...]: def _nested_plugin_value(root: object, segments: tuple[str, ...], default: Any) -> Any: + """Walk ``segments`` through nested mappings; ``default`` on the first miss.""" current = root for segment in segments: if not isinstance(current, Mapping) or segment not in current: @@ -61,12 +50,19 @@ def _nested_plugin_value(root: object, segments: tuple[str, ...], default: Any) def _nested_plugin_mapping(segments: tuple[str, ...], value: Any) -> dict[str, Any]: + """Wrap ``value`` in nested single-key dicts, outermost first.""" nested: Any = value for segment in reversed(segments): nested = {segment: nested} return nested +def _plugin_settings_entry(config: object, plugin_id: str) -> Mapping[str, Any] | None: + """``plugins.entries.`` as a mapping, else ``None``.""" + entry = _nested_plugin_value(config, ("plugins", "entries", plugin_id), None) + return entry if isinstance(entry, Mapping) else None + + def _plugin_data_namespace(plugin_id: str, skill_namespace: str) -> str: """Return one Windows-safe directory component for plugin-owned data.""" candidate = skill_namespace or plugin_id @@ -80,24 +76,20 @@ def _plugin_data_namespace(plugin_id: str, skill_namespace: str) -> str: return _portable_skill_namespace(candidate) -def _state_thread_lock(path: Path) -> threading.RLock: - key = str(path.resolve(strict=False)) - with _PLUGIN_STATE_LOCKS_GUARD: - return _PLUGIN_STATE_LOCKS.setdefault(key, threading.RLock()) - - @contextmanager def _locked_plugin_state(path: Path): """Serialize state read-modify-write across threads/processes (fcntl / msvcrt). The lock lives in a sibling file because atomic replacement changes the target's inode.""" lock_path = path.with_name(f".{path.name}.lock") - thread_lock = _state_thread_lock(lock_path) + with _PLUGIN_STATE_LOCKS_GUARD: + thread_lock = _PLUGIN_STATE_LOCKS.setdefault( + str(lock_path.resolve(strict=False)), threading.RLock() + ) with thread_lock: lock_path.parent.mkdir(parents=True, exist_ok=True) with open(lock_path, "a+b") as handle: if os.name == "nt": # pragma: no cover - exercised on Windows CI import msvcrt - if handle.seek(0, os.SEEK_END) == 0: handle.write(b"\0") handle.flush() @@ -105,7 +97,6 @@ def _locked_plugin_state(path: Path): msvcrt.locking(handle.fileno(), msvcrt.LK_LOCK, 1) else: import fcntl - fcntl.flock(handle.fileno(), fcntl.LOCK_EX) try: yield @@ -138,7 +129,7 @@ class PluginState: @staticmethod def _validate_key(key: str) -> None: - if (not isinstance(key, str) or not _PLUGIN_STATE_KEY_RE.fullmatch(key) or ".." in key): + if not isinstance(key, str) or not _PLUGIN_STATE_KEY_RE.fullmatch(key) or ".." in key: raise ValueError( "Plugin state keys must be 1-128 characters using letters, " "numbers, '_', '-', '.', or ':' (without '..')" @@ -180,5 +171,4 @@ class PluginState: f"than the {self.quota_bytes}-byte per-plugin quota" ) from utils import atomic_json_write - atomic_json_write(self.path, data, mode=0o600)