mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-07-31 19:16:29 +00:00
feat(nemo-relay): activate dynamic plugins
Signed-off-by: Bryan Bednarski <bbednarski@nvidia.com>
This commit is contained in:
parent
7426c09bee
commit
0c8cf21882
4 changed files with 388 additions and 22 deletions
|
|
@ -84,14 +84,15 @@ wheel from this checkout, then install the official NeMo Relay runtime extra:
|
|||
```bash
|
||||
uv build --wheel
|
||||
python -m pip install --force-reinstall dist/hermes_agent-*.whl
|
||||
python -m pip install "nemo-relay==0.3"
|
||||
python -m pip install "nemo-relay>=0.5,<0.7"
|
||||
hermes plugins enable observability/nemo_relay
|
||||
```
|
||||
|
||||
The plugin fails open when `nemo-relay` is not installed. Install and test it against the official NeMo Relay 0.3 PyPI distribution:
|
||||
The plugin fails open when `nemo-relay` is not installed. Install a supported
|
||||
NeMo Relay 0.5 or 0.6 distribution:
|
||||
|
||||
```bash
|
||||
pip install "nemo-relay==0.3"
|
||||
pip install "nemo-relay>=0.5,<0.7"
|
||||
```
|
||||
|
||||
## Export Configuration
|
||||
|
|
@ -189,16 +190,46 @@ session, turn, approval, and subagent marks; the plugin skips its manual
|
|||
NeMo Relay. `tool_parallelism.mode = "observe_only"` keeps tool scheduling
|
||||
observational while still wrapping the real execution boundary.
|
||||
|
||||
### Dynamic Plugins (NeMo Relay 0.6)
|
||||
|
||||
Hermes can activate explicit native or worker plugins through the NeMo Relay
|
||||
0.6 binding. Add `dynamic_plugins` entries to the same file; each entry uses
|
||||
the binding's `plugin_id`, `kind`, `manifest_ref`, optional `environment_ref`,
|
||||
and component `config` fields:
|
||||
|
||||
```toml
|
||||
[[dynamic_plugins]]
|
||||
plugin_id = "example-plugin"
|
||||
kind = "rust_dynamic"
|
||||
manifest_ref = "/absolute/path/to/example-plugin/relay-plugin.toml"
|
||||
|
||||
[dynamic_plugins.config]
|
||||
mode = "enabled"
|
||||
```
|
||||
|
||||
Hermes activates these plugins before registering its managed LLM and tool
|
||||
execution middleware and retains the activation for the runtime lifetime.
|
||||
During shutdown it closes session exporters, flushes Relay subscribers, and
|
||||
then closes the activation so callbacks are removed before plugin code is
|
||||
unloaded.
|
||||
|
||||
NeMo Relay 0.5 does not expose dynamic activation through its Python binding.
|
||||
When `dynamic_plugins` is present with a 0.5 runtime, Hermes logs an actionable
|
||||
warning and continues with the ordinary static component configuration, so
|
||||
ATOF and ATIF observability remain available. No dynamic plugin is loaded in
|
||||
that degraded mode.
|
||||
|
||||
For the full generic Hermes middleware contract, see
|
||||
[`docs/middleware/README.md`](../../../docs/middleware/README.md).
|
||||
|
||||
## Canonical Local Examples
|
||||
|
||||
The observe-only examples in this section use the official `nemo-relay==0.3`
|
||||
distribution and a local Ollama model served through the OpenAI-compatible API.
|
||||
The observe-only examples in this section use a supported NeMo Relay 0.5 or
|
||||
0.6 distribution and a local Ollama model served through the OpenAI-compatible
|
||||
API.
|
||||
|
||||
```bash
|
||||
pip install "nemo-relay==0.3"
|
||||
pip install "nemo-relay>=0.5,<0.7"
|
||||
|
||||
export HERMES_HOME=/tmp/hermes-nemo-relay-docs/hermes-home
|
||||
mkdir -p "$HERMES_HOME"
|
||||
|
|
@ -444,9 +475,8 @@ for the same execution.
|
|||
|
||||
This example enables both NeMo Relay observability export and adaptive execution
|
||||
middleware for a local Hermes run. This path requires a NeMo Relay runtime that
|
||||
supports `[components.config.tool_parallelism]`; the `nemo-relay==0.3`
|
||||
install used by the earlier observability-only examples does not support this
|
||||
adaptive config.
|
||||
supports `[components.config.tool_parallelism]`, as provided by the supported
|
||||
0.5 and 0.6 releases.
|
||||
|
||||
```bash
|
||||
export HERMES_HOME=/tmp/hermes-middleware-test/hermes-home
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
import atexit
|
||||
import asyncio
|
||||
import inspect
|
||||
import json
|
||||
|
|
@ -44,6 +45,7 @@ class _SubagentParent:
|
|||
class _Settings:
|
||||
plugins_toml_path: str = ""
|
||||
plugins_config: dict[str, Any] | None = None
|
||||
dynamic_plugins: list[dict[str, Any]] = field(default_factory=list)
|
||||
adaptive_enabled: bool = False
|
||||
adaptive_mode: str = "observe_only"
|
||||
atof_enabled: bool = False
|
||||
|
|
@ -67,6 +69,8 @@ class _Runtime:
|
|||
self.subagent_parents: dict[str, _SubagentParent] = {}
|
||||
self.atof_exporter: Any = None
|
||||
self._atof_subscriber_name = "hermes.nemo_relay.atof"
|
||||
self._plugin_activation: Any = None
|
||||
self._shutdown_registered = False
|
||||
self._plugin_config_initialized = self._configure_plugins_toml()
|
||||
self._plugin_config_needs_reinit = False
|
||||
if not self._plugin_config_initialized:
|
||||
|
|
@ -76,27 +80,64 @@ class _Runtime:
|
|||
if not self.settings.plugins_config:
|
||||
return False
|
||||
plugin_mod = getattr(self.nemo_relay, "plugin", None)
|
||||
if plugin_mod is None:
|
||||
return False
|
||||
plugin_config = _static_plugin_config(self.settings.plugins_config)
|
||||
if self.settings.dynamic_plugins:
|
||||
activate_dynamic = getattr(plugin_mod, "activate_dynamic_plugins", None)
|
||||
if callable(activate_dynamic):
|
||||
try:
|
||||
self._ensure_plugin_config_output_dirs(plugin_config)
|
||||
self._plugin_activation = _resolve_awaitable(
|
||||
activate_dynamic(plugin_config, self.settings.dynamic_plugins)
|
||||
)
|
||||
self._ensure_shutdown_registered()
|
||||
return True
|
||||
except Exception as exc:
|
||||
logger.warning(
|
||||
"NeMo Relay dynamic plugin activation failed; continuing with static "
|
||||
"observability only: %s",
|
||||
exc,
|
||||
)
|
||||
else:
|
||||
logger.warning(
|
||||
"NeMo Relay dynamic plugins require nemo-relay>=0.6,<0.7; this binding does "
|
||||
"not expose plugin.activate_dynamic_plugins. Continuing with static "
|
||||
"observability only."
|
||||
)
|
||||
initialize = getattr(plugin_mod, "initialize", None)
|
||||
if not callable(initialize):
|
||||
return False
|
||||
try:
|
||||
self._ensure_plugin_config_output_dirs(self.settings.plugins_config)
|
||||
_resolve_awaitable(initialize(self.settings.plugins_config))
|
||||
self._ensure_plugin_config_output_dirs(plugin_config)
|
||||
_resolve_awaitable(initialize(plugin_config))
|
||||
return True
|
||||
except Exception as exc:
|
||||
logger.debug("NeMo Relay plugins.toml init failed: %s", exc, exc_info=True)
|
||||
return False
|
||||
|
||||
def _ensure_shutdown_registered(self) -> None:
|
||||
if self._shutdown_registered:
|
||||
return
|
||||
atexit.register(self.shutdown)
|
||||
self._shutdown_registered = True
|
||||
|
||||
def _clear_plugins_toml(self) -> None:
|
||||
if not self._plugin_config_initialized:
|
||||
return
|
||||
plugin_mod = getattr(self.nemo_relay, "plugin", None)
|
||||
clear = getattr(plugin_mod, "clear", None)
|
||||
if not callable(clear):
|
||||
return
|
||||
try:
|
||||
_resolve_awaitable(clear())
|
||||
if self._plugin_activation is not None:
|
||||
_flush_relay_subscribers(self.nemo_relay)
|
||||
close = getattr(self._plugin_activation, "close", None)
|
||||
if callable(close):
|
||||
_resolve_awaitable(close())
|
||||
else:
|
||||
plugin_mod = getattr(self.nemo_relay, "plugin", None)
|
||||
clear = getattr(plugin_mod, "clear", None)
|
||||
if callable(clear):
|
||||
_resolve_awaitable(clear())
|
||||
finally:
|
||||
self._plugin_activation = None
|
||||
self._plugin_config_initialized = False
|
||||
self._plugin_config_needs_reinit = bool(self.settings.plugins_config)
|
||||
|
||||
|
|
@ -226,20 +267,46 @@ class _Runtime:
|
|||
self.nemo_relay.scope.pop(state.handle, output=_jsonable(kwargs))
|
||||
except Exception:
|
||||
logger.debug("NeMo Relay session pop failed", exc_info=True)
|
||||
try:
|
||||
_flush_relay_subscribers(self.nemo_relay)
|
||||
except Exception:
|
||||
logger.debug("NeMo Relay subscriber flush failed", exc_info=True)
|
||||
self.export_atif(state)
|
||||
if state.atif_exporter is not None and state.atif_subscriber_name:
|
||||
try:
|
||||
state.atif_exporter.deregister(state.atif_subscriber_name)
|
||||
except Exception:
|
||||
logger.debug("NeMo Relay ATIF deregister failed", exc_info=True)
|
||||
if self._plugin_config_initialized and not self.sessions:
|
||||
if (
|
||||
self._plugin_config_initialized
|
||||
and self._plugin_activation is None
|
||||
and not self.sessions
|
||||
):
|
||||
try:
|
||||
self._clear_plugins_toml()
|
||||
except Exception:
|
||||
logger.debug("NeMo Relay plugins.toml clear failed", exc_info=True)
|
||||
elif self.settings.plugins_config and not self.sessions:
|
||||
elif (
|
||||
self.settings.plugins_config
|
||||
and self._plugin_activation is None
|
||||
and not self.sessions
|
||||
):
|
||||
self._plugin_config_needs_reinit = True
|
||||
|
||||
def shutdown(self) -> None:
|
||||
"""Close active sessions and the process-lifetime plugin activation."""
|
||||
for session_id in list(self.sessions):
|
||||
self.close_session({"session_id": session_id, "reason": "runtime_shutdown"})
|
||||
if self._plugin_config_initialized:
|
||||
try:
|
||||
self._clear_plugins_toml()
|
||||
except Exception:
|
||||
logger.debug("NeMo Relay plugin runtime shutdown failed", exc_info=True)
|
||||
self._clear_atof()
|
||||
if self._shutdown_registered:
|
||||
atexit.unregister(self.shutdown)
|
||||
self._shutdown_registered = False
|
||||
|
||||
def mark(self, name: str, kwargs: dict[str, Any]) -> None:
|
||||
state = self.ensure_session(kwargs)
|
||||
self.nemo_relay.scope.event(
|
||||
|
|
@ -274,14 +341,14 @@ class _Runtime:
|
|||
|
||||
def managed_llm_enabled(self) -> bool:
|
||||
return (
|
||||
self.settings.adaptive_enabled
|
||||
(self.settings.adaptive_enabled or self._plugin_activation is not None)
|
||||
and callable(getattr(getattr(self.nemo_relay, "llm", None), "execute", None))
|
||||
and callable(getattr(self.nemo_relay, "LLMRequest", None))
|
||||
)
|
||||
|
||||
def managed_tool_enabled(self) -> bool:
|
||||
return (
|
||||
self.settings.adaptive_enabled
|
||||
(self.settings.adaptive_enabled or self._plugin_activation is not None)
|
||||
and callable(getattr(getattr(self.nemo_relay, "tools", None), "execute", None))
|
||||
)
|
||||
|
||||
|
|
@ -402,6 +469,10 @@ class _Runtime:
|
|||
|
||||
|
||||
def register(ctx) -> None:
|
||||
# Activate dynamic plugins before Hermes installs the managed execution
|
||||
# boundaries that invoke their interceptors.
|
||||
if _load_settings().dynamic_plugins:
|
||||
_get_runtime()
|
||||
ctx.register_hook("on_session_start", on_session_start)
|
||||
ctx.register_hook("on_session_end", on_session_end)
|
||||
ctx.register_hook("on_session_finalize", on_session_finalize)
|
||||
|
|
@ -648,6 +719,7 @@ def _load_settings() -> _Settings:
|
|||
return _Settings(
|
||||
plugins_toml_path=plugins_toml_path,
|
||||
plugins_config=plugins_config,
|
||||
dynamic_plugins=_dynamic_plugin_specs(plugins_config),
|
||||
adaptive_enabled=adaptive_config is not None,
|
||||
adaptive_mode=_adaptive_mode(adaptive_config),
|
||||
atof_enabled=_env_bool("HERMES_NEMO_RELAY_ATOF_ENABLED"),
|
||||
|
|
@ -664,6 +736,81 @@ def _load_settings() -> _Settings:
|
|||
)
|
||||
|
||||
|
||||
def _static_plugin_config(plugins_config: dict[str, Any]) -> dict[str, Any]:
|
||||
"""Return Relay's base plugin config without Hermes-owned host fields."""
|
||||
return {key: value for key, value in plugins_config.items() if key != "dynamic_plugins"}
|
||||
|
||||
|
||||
def _dynamic_plugin_specs(plugins_config: dict[str, Any] | None) -> list[dict[str, Any]]:
|
||||
if not isinstance(plugins_config, dict):
|
||||
return []
|
||||
raw_specs = plugins_config.get("dynamic_plugins")
|
||||
if raw_specs is None:
|
||||
return []
|
||||
if not isinstance(raw_specs, list):
|
||||
logger.warning(
|
||||
"Ignoring invalid NeMo Relay dynamic_plugins config: expected an array of plugin specs"
|
||||
)
|
||||
return []
|
||||
|
||||
specs: list[dict[str, Any]] = []
|
||||
for index, raw_spec in enumerate(raw_specs):
|
||||
if not isinstance(raw_spec, dict):
|
||||
logger.warning(
|
||||
"Ignoring invalid NeMo Relay dynamic_plugins[%d]: expected an object", index
|
||||
)
|
||||
continue
|
||||
plugin_id = raw_spec.get("plugin_id")
|
||||
kind = raw_spec.get("kind")
|
||||
manifest_ref = raw_spec.get("manifest_ref")
|
||||
config = raw_spec.get("config", {})
|
||||
environment_ref = raw_spec.get("environment_ref")
|
||||
if not isinstance(plugin_id, str) or not plugin_id.strip():
|
||||
logger.warning(
|
||||
"Ignoring invalid NeMo Relay dynamic_plugins[%d]: plugin_id is required", index
|
||||
)
|
||||
continue
|
||||
if kind not in {"rust_dynamic", "worker"}:
|
||||
logger.warning(
|
||||
"Ignoring invalid NeMo Relay dynamic_plugins[%d]: kind must be rust_dynamic or worker",
|
||||
index,
|
||||
)
|
||||
continue
|
||||
if not isinstance(manifest_ref, str) or not manifest_ref.strip():
|
||||
logger.warning(
|
||||
"Ignoring invalid NeMo Relay dynamic_plugins[%d]: manifest_ref is required", index
|
||||
)
|
||||
continue
|
||||
if not isinstance(config, dict):
|
||||
logger.warning(
|
||||
"Ignoring invalid NeMo Relay dynamic_plugins[%d]: config must be an object", index
|
||||
)
|
||||
continue
|
||||
if environment_ref is not None and not isinstance(environment_ref, str):
|
||||
logger.warning(
|
||||
"Ignoring invalid NeMo Relay dynamic_plugins[%d]: environment_ref must be a string",
|
||||
index,
|
||||
)
|
||||
continue
|
||||
spec: dict[str, Any] = {
|
||||
"plugin_id": plugin_id,
|
||||
"kind": kind,
|
||||
"manifest_ref": manifest_ref,
|
||||
"config": config,
|
||||
}
|
||||
if environment_ref is not None:
|
||||
spec["environment_ref"] = environment_ref
|
||||
specs.append(spec)
|
||||
return specs
|
||||
|
||||
|
||||
def _flush_relay_subscribers(nemo_relay: Any) -> None:
|
||||
subscribers = getattr(nemo_relay, "subscribers", None)
|
||||
flush = getattr(subscribers, "flush", None)
|
||||
if callable(flush):
|
||||
flush()
|
||||
|
||||
|
||||
def _load_plugins_config(path: str) -> dict[str, Any] | None:
|
||||
if not path:
|
||||
return None
|
||||
|
|
@ -959,4 +1106,6 @@ def _resolve_awaitable(value: Any) -> Any:
|
|||
def reset_for_tests() -> None:
|
||||
global _RUNTIME
|
||||
with _LOCK:
|
||||
if isinstance(_RUNTIME, _Runtime):
|
||||
_RUNTIME.shutdown()
|
||||
_RUNTIME = None
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue