From 8afd98ef2af3988ba7fe8a9a29e07d53ca4cb42b Mon Sep 17 00:00:00 2001 From: Bryan Bednarski Date: Mon, 10 Aug 2026 18:19:40 -0600 Subject: [PATCH] refactor(relay): remove legacy observability plugin Signed-off-by: Bryan Bednarski --- docs/middleware/README.md | 5 +- docs/observability/README.md | 8 +- docs/observability/monitoring.md | 4 +- docs/observability/relay-shared-metrics.md | 19 +- hermes_cli/plugins_cmd.py | 10 +- plugins/observability/nemo_relay/README.md | 602 ---------- plugins/observability/nemo_relay/__init__.py | 1023 ----------------- plugins/observability/nemo_relay/plugin.yaml | 15 - scripts/toolperf_abeval/README.md | 4 +- scripts/toolperf_abeval/ab_eval.py | 28 +- .../test_plugins_cmd_enable_disable_nested.py | 14 +- tests/plugins/test_nemo_relay_plugin.py | 490 -------- .../user-guide/features/built-in-plugins.md | 1 - 13 files changed, 59 insertions(+), 2164 deletions(-) delete mode 100644 plugins/observability/nemo_relay/README.md delete mode 100644 plugins/observability/nemo_relay/__init__.py delete mode 100644 plugins/observability/nemo_relay/plugin.yaml delete mode 100644 tests/plugins/test_nemo_relay_plugin.py diff --git a/docs/middleware/README.md b/docs/middleware/README.md index 4a5c06f8cb..96304fa0e1 100644 --- a/docs/middleware/README.md +++ b/docs/middleware/README.md @@ -233,8 +233,9 @@ Execution middleware may call `next_call(modified_args)` to pass a changed payload to later middleware and the base tool dispatcher. Plugin-specific examples should live with the plugin that owns the behavior. -For NeMo Relay adaptive execution middleware, see -[`plugins/observability/nemo_relay/README.md`](../../plugins/observability/nemo_relay/README.md). +NeMo Relay execution middleware is installed through an explicitly selected +Relay `plugins.toml`; see +[Relay shared metrics](../observability/relay-shared-metrics.md). ## Safety Notes diff --git a/docs/observability/README.md b/docs/observability/README.md index 915b0d5e2a..f5ebea6177 100644 --- a/docs/observability/README.md +++ b/docs/observability/README.md @@ -318,7 +318,7 @@ nested agent work or security lifecycle events. The bundled Langfuse plugin demonstrates direct hook-based observability for turns, provider requests, and tool calls. -The bundled NeMo Relay plugin maps the same generic observer contract to NeMo -Relay scopes, LLM spans, tool spans, marks, ATOF streams, and ATIF exports. -NeMo Relay-specific configuration and examples live in -[`plugins/observability/nemo_relay/README.md`](../../plugins/observability/nemo_relay/README.md). +The native NeMo Relay SDK integration maps Hermes session, turn, LLM, tool, +and mark lifecycles to Relay. Explicit Relay plugin configuration can add ATOF +or ATIF exporters and execution middleware; see +[Relay shared metrics](relay-shared-metrics.md). diff --git a/docs/observability/monitoring.md b/docs/observability/monitoring.md index f798bf2dbe..cba4ead1ff 100644 --- a/docs/observability/monitoring.md +++ b/docs/observability/monitoring.md @@ -9,8 +9,8 @@ lifecycle state, platform connector health, and content-free warning/error diagnostics. It never exports prompts, messages, tool arguments or results, job names, destinations, schedules, raw errors, session history, usage analytics, audit logs, or detailed execution traces. Run/model/tool trajectory -capture is a separate plane served by the NeMo Relay integration -(`plugins/observability/nemo_relay/`) and its Hermes-owned subscribers. +capture is a separate plane served by Hermes's native NeMo Relay SDK +integration and explicitly configured Relay subscribers or exporters. ## What gets exported diff --git a/docs/observability/relay-shared-metrics.md b/docs/observability/relay-shared-metrics.md index 860d916e21..3ae597b5d6 100644 --- a/docs/observability/relay-shared-metrics.md +++ b/docs/observability/relay-shared-metrics.md @@ -2,21 +2,21 @@ Hermes includes NeMo Relay as a normal runtime dependency on platforms for which Relay publishes a native wheel. The shared-metrics integration is built -into Hermes and does not require `hermes plugins enable -observability/nemo_relay`. Hermes remains importable without Relay on other -native targets. Those targets use an explicit reduced-capability no-op host: +into Hermes and does not require a Hermes observability plugin. Hermes remains +importable without Relay on other native targets. Those targets use an +explicit reduced-capability no-op host: Hermes execution remains available, while Relay scopes, middleware, plugins, and subscribers are unavailable. The `hermes-agent[nemo-relay]` extra remains as a no-op compatibility alias for existing installation commands. -Hermes requires NeMo Relay 0.6.0 or later within the 0.6 release line. That +Hermes requires NeMo Relay 0.7.1 or later within the 0.7 release line. That release establishes the lossless provider-codec contract used for Anthropic Messages, OpenAI Chat Completions, and OpenAI Responses requests. ## Runtime Dependency and Data Boundary Hermes installs the platform-specific `nemo-relay` native wheel from the -bounded `>=0.6.0,<0.7` dependency range. The published package is built from +bounded `>=0.7.1,<0.8` dependency range. The published package is built from the [NVIDIA NeMo Relay repository](https://github.com/NVIDIA/NeMo-Relay). Unsupported platforms use the explicit no-op runtime described above rather than downloading a different implementation. @@ -41,9 +41,12 @@ This choice is read from the profile's own `config.yaml`. A machine-managed configuration overlay cannot enable or disable shared metrics on the profile's behalf. -The existing `observability/nemo_relay` plugin remains separate. Enable that -plugin only for its opt-in rich observability exporters, adaptive execution, -or dynamic Relay plugins. +Relay plugin activation is owned by the native runtime and remains explicitly +opt-in. Set `HERMES_NEMO_RELAY_PLUGINS_TOML` to a selected `plugins.toml` to +activate configured middleware, exporters, or dynamic plugins. When it is +unset, Hermes does not invoke Relay's plugin initializer or trigger Relay +plugin configuration discovery. Invalid explicit configuration is reported +and Hermes continues without native plugin activation. Hermes core owns one Relay host and one isolated Relay session scope per Hermes session. Core lifecycle producers use diff --git a/hermes_cli/plugins_cmd.py b/hermes_cli/plugins_cmd.py index 6f7446620e..2c258ee8e8 100644 --- a/hermes_cli/plugins_cmd.py +++ b/hermes_cli/plugins_cmd.py @@ -1336,14 +1336,14 @@ def _save_enabled_set(enabled: set) -> None: def _resolve_plugin_key(name: str) -> Optional[str]: """Resolve a user-supplied plugin identifier to its canonical registry key. - Accepts either the bare manifest name (``nemo_relay``), the directory - name, or the full path-derived key (``observability/nemo_relay``) and + Accepts either the bare manifest name (``langfuse``), the directory + name, or the full path-derived key (``observability/langfuse``) and returns the canonical key the loader gates on (``manifest.key`` or, for a flat plugin, the bare name). Returns ``None`` when no plugin matches. This is the single normalization point so ``hermes plugins enable`` / ``disable`` write the same key that ``PluginManager`` matches against — - nested category plugins (e.g. ``observability/nemo_relay``) included. + nested category plugins (e.g. ``observability/langfuse``) included. """ entries = _discover_all_plugins() # 1. Exact match on canonical key or manifest name — always unambiguous. @@ -1351,8 +1351,8 @@ def _resolve_plugin_key(name: str) -> Optional[str]: # entry = (name, version, description, source, dir_path, key) if name == entry[5] or name == entry[0]: return entry[5] - # 2. Fall back to a bare leaf-name match (e.g. "nemo_relay" -> - # "observability/nemo_relay"), but only when it resolves to exactly one + # 2. Fall back to a bare leaf-name match (e.g. "langfuse" -> + # "observability/langfuse"), but only when it resolves to exactly one # plugin so we never silently pick the wrong same-named nested plugin. leaf_matches = [entry[5] for entry in entries if name == entry[5].split("/")[-1]] if len(leaf_matches) == 1: diff --git a/plugins/observability/nemo_relay/README.md b/plugins/observability/nemo_relay/README.md deleted file mode 100644 index 85f7d03903..0000000000 --- a/plugins/observability/nemo_relay/README.md +++ /dev/null @@ -1,602 +0,0 @@ -# NeMo Relay Observability - -Optional Hermes observability plugin that configures exporters and maps -Hermes-specific observer hooks to NeMo Relay marks and ATIF state. Hermes core -owns Relay session, turn, LLM, and tool execution scopes. - -NeMo Relay is NVIDIA's runtime layer for agent execution boundaries. It does -not replace Hermes Agent's planner, tools, memory, model provider routing, or -CLI UX. Hermes core emits NeMo Relay lifecycle events for provider and tool -execution, while this plugin enables rich exporters and observer marks for -sessions, turns, approval prompts, and delegated subagents. - -With this plugin enabled, Hermes Agent can: - -- Export the Relay scopes and LLM/tool lifecycles emitted by Hermes core. -- Add Hermes session, turn, approval, and subagent mark events. -- Export raw lifecycle events as Agent Trajectory Observability Format (ATOF) - JSONL for debugging and offline inspection. -- Export Agent Trajectory Interchange Format (ATIF) trajectories for replay, - evaluation, and harness analysis workflows. -- Correlate parent sessions, delegated subagents, tool calls, and provider - calls through shared session, turn, and trajectory metadata. - -See the NeMo Relay overview for the broader runtime model: -https://docs.nvidia.com/nemo/relay/about-nemo-relay/overview - -ATOF is NVIDIA's canonical JSONL event stream representation for NeMo Relay -lifecycle events. The format is documented in the NeMo Agent Toolkit: -https://github.com/NVIDIA/NeMo-Agent-Toolkit/blob/develop/packages/nvidia_nat_atif/atof-event-format.md - -ATIF is the trajectory representation produced from those events. NVIDIA and -Harbor upstreamed ATIF v1.7 support for complex harness workflows, including -subagent trajectory embedding, trajectory IDs, multi-LLM-call step metadata, and -deterministic no-LLM orchestration steps: -https://github.com/harbor-framework/harbor/blob/main/rfcs/0001-trajectory-format.md - -## Enablement - -Enable the plugin before setting export options: - -```bash -hermes plugins enable observability/nemo_relay -``` - -The `HERMES_NEMO_RELAY_*` environment variables below only configure an -already-enabled plugin. They do not enable plugin discovery by themselves. - -For isolated test homes, enable the plugin in the same `HERMES_HOME` that the -agent run will use: - -```bash -env HERMES_HOME=/tmp/hermes-nemo-relay-test \ - hermes plugins enable observability/nemo_relay -``` - -Runs started with `--ignore_user_config` skip the enabled-plugin state from -`HERMES_HOME`, so local E2E tests should omit that flag unless the test harness -loads `observability/nemo_relay` explicitly another way. - -`HERMES_HOME` is the Hermes profile/config home used by both -`hermes plugins enable ...` and the later `hermes chat ...` run. If unset, -Hermes uses the user's default home, usually `~/.hermes`. For isolated smoke -tests, choose any writable temporary directory and use the same value for every -command in that test: - -```bash -export HERMES_HOME=/tmp/hermes-nemo-relay-test -hermes plugins enable observability/nemo_relay -hermes chat --query 'Reply exactly ok' --provider custom --model qwen3.6:35b -``` - -For source checkouts, make sure the `hermes` command you run is built from the -checkout that contains this plugin. A globally installed older CLI will not see -new bundled plugins from your working tree. - -```bash -uv sync -uv run hermes plugins enable observability/nemo_relay -uv run hermes chat --query 'Reply exactly ok' --provider custom --model qwen3.6:35b -``` - -To ship the updated CLI into another environment, build and install a fresh -wheel from this checkout. On platforms for which Relay publishes a native -wheel, Hermes installs its supported NeMo Relay runtime as a normal dependency: - -```bash -uv build --wheel -python -m pip install --force-reinstall dist/hermes_agent-*.whl -hermes plugins enable observability/nemo_relay -``` - -The plugin remains opt-in even though the runtime dependency is installed by -default. Enabling this plugin controls rich observability and adaptive -behavior; it does not control Hermes shared client metrics. - -## Export Configuration - -The plugin can configure exporters directly from `HERMES_NEMO_RELAY_*` -environment variables, or delegate exporter setup to a NeMo Relay -`plugins.toml` component config. - -Use environment variables for local smoke tests, CI jobs, and one-off CLI -runs. Use `plugins.toml` when you want one NeMo Relay configuration document to -own observability components such as ATOF, ATIF, OpenTelemetry, and -OpenInference. - -### Environment Variables - -Useful local export settings after the plugin is enabled: - -```bash -export HERMES_NEMO_RELAY_ATOF_ENABLED=1 -export HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY=.nemo-relay/atof -export HERMES_NEMO_RELAY_ATIF_ENABLED=1 -export HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY=.nemo-relay/atif -``` - -Optional overrides: - -- `HERMES_NEMO_RELAY_ATOF_FILENAME` -- `HERMES_NEMO_RELAY_ATOF_MODE` (`append` or `overwrite`) -- `HERMES_NEMO_RELAY_ATIF_FILENAME_TEMPLATE` -- `HERMES_NEMO_RELAY_ATIF_AGENT_NAME` -- `HERMES_NEMO_RELAY_ATIF_AGENT_VERSION` -- `HERMES_NEMO_RELAY_ATIF_MODEL_NAME` -- `HERMES_NEMO_RELAY_ATIF_SUBAGENT_EXPORT_MODE` (`embedded` by default; set `all` to also write standalone child files) - -### NeMo Relay Component Config - -To initialize NeMo Relay from a component config, create a `plugins.toml` file -and point Hermes at it: - -```bash -export HERMES_NEMO_RELAY_PLUGINS_TOML=.nemo-relay/plugins.toml -``` - -Minimal ATOF and ATIF config: - -```toml -version = 1 - -[[components]] -kind = "observability" -enabled = true - -[components.config] -version = 1 - -[components.config.atof] -enabled = true -output_directory = ".nemo-relay/atof" -filename = "events.jsonl" -mode = "overwrite" - -[components.config.atif] -enabled = true -output_directory = ".nemo-relay/atif" -filename_template = "trajectory-{session_id}.json" -agent_name = "Hermes Agent" -agent_version = "local" -``` - -When `HERMES_NEMO_RELAY_PLUGINS_TOML` is set and initializes successfully, NeMo -Relay owns exporter lifecycle through that config. The direct -`HERMES_NEMO_RELAY_ATOF_*` fallback setup is skipped. If the same -`plugins.toml` observability config enables `atif`, the direct -`HERMES_NEMO_RELAY_ATIF_*` fallback setup is also skipped so Hermes does not -double-export trajectories on teardown. If `plugins.toml` initialization fails, -Hermes keeps the direct env-var fallbacks active for that run. - -Hermes core routes provider and tool execution through NeMo Relay managed APIs -regardless of whether this plugin is enabled. To install adaptive interceptors -on those boundaries, include an adaptive component in the same `plugins.toml`: - -```toml -[[components]] -kind = "adaptive" -enabled = true - -[components.config.tool_parallelism] -mode = "observe_only" -``` - -The observer hooks emit session, turn, approval, and subagent marks. They do not -create a second LLM or tool lifecycle. `tool_parallelism.mode = "observe_only"` -keeps tool scheduling observational while still intercepting the core-managed -execution boundary. - -### Dynamic Plugins - -Hermes uses the dynamic-plugin activation API available in NeMo Relay 0.6 and -later. Configure native or worker plugins with Hermes-owned -`[[dynamic_plugins]]` entries that match the Python binding's activation-spec -fields: - -```toml -[[dynamic_plugins]] -plugin_id = "example-plugin" -kind = "rust_dynamic" -manifest_ref = "./example-plugin/relay-plugin.toml" - -[dynamic_plugins.config] -mode = "enabled" -``` - -For a worker plugin, also provide the lifecycle-managed `environment_ref`: - -```toml -[[dynamic_plugins]] -plugin_id = "example-worker" -kind = "worker" -manifest_ref = "./example-worker/relay-plugin.toml" -environment_ref = "/absolute/path/from-nemo-relay-plugins-inspect" - -[dynamic_plugins.config] -mode = "enabled" -``` - -Provision the worker first with `nemo-relay plugins add`, then copy -`data.source.environment_ref` from the JSON output of -`nemo-relay plugins inspect --json`. Relay rejects arbitrary Python -environments at activation time. - -Relative `manifest_ref` and `environment_ref` values resolve relative to the -physical `plugins.toml` file. - -Relay's canonical gateway `[[plugins.dynamic]]` records are not interchangeable -with this Hermes-owned section. The gateway combines those records with -separate lifecycle state for enablement, trust policy, and worker environments; -the Python binding does not yet expose that resolver. Hermes rejects -`[[plugins.dynamic]]` with an actionable diagnostic instead of silently -ignoring it or bypassing lifecycle policy. Use `[[dynamic_plugins]]` until Relay -exposes shared file-and-lifecycle resolution to embedding hosts. - -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. - -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 NeMo Relay runtime installed -with Hermes and a local Ollama model served through the OpenAI-compatible API. - -```bash -export HERMES_HOME=/tmp/hermes-nemo-relay-docs/hermes-home -mkdir -p "$HERMES_HOME" - -cat > "$HERMES_HOME/config.yaml" <<'YAML' -model: - provider: custom - default: qwen3.6:35b - base_url: http://127.0.0.1:11434/v1 - api_key: ollama -plugins: - enabled: - - observability/nemo_relay -delegation: - max_spawn_depth: 2 - max_concurrent_children: 2 - child_timeout_seconds: 180 - model: qwen3.6:35b - provider: custom - base_url: http://127.0.0.1:11434/v1 - api_key: ollama -YAML -``` - -### Delegated Subagent Tool Call - -This run starts a parent Hermes session, delegates to a child subagent, has the -child call `terminal`, and writes both ATOF and ATIF. - -```bash -export HERMES_NEMO_RELAY_ATOF_ENABLED=1 -export HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY=/tmp/hermes-nemo-relay-docs/subagent/atof -export HERMES_NEMO_RELAY_ATOF_FILENAME=nested-subagent-atof.jsonl -export HERMES_NEMO_RELAY_ATOF_MODE=overwrite -export HERMES_NEMO_RELAY_ATIF_ENABLED=1 -export HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY=/tmp/hermes-nemo-relay-docs/subagent/atif -export HERMES_NEMO_RELAY_ATIF_FILENAME_TEMPLATE='nested-subagent-atif-{session_id}.json' -export HERMES_NEMO_RELAY_ATIF_AGENT_NAME='Hermes Agent E2E' -export HERMES_NEMO_RELAY_ATIF_AGENT_VERSION=docs-example -export HERMES_NEMO_RELAY_ATIF_SUBAGENT_EXPORT_MODE=all - -hermes chat \ - --query 'Use delegate_task exactly once. Ask the child subagent to use the terminal tool exactly once to run printf docs_nested_leaf_function. After the child returns, reply with exactly: parent received nested subagent result.' \ - --provider custom \ - --model qwen3.6:35b \ - --toolsets delegation,terminal \ - --max-turns 10 \ - --quiet \ - --accept-hooks -``` - -CLI output: - -```text -session_id: docs-parent-session -parent received nested subagent result. -``` - -Sanitized ATOF excerpt: - -```jsonl -{"kind":"scope","category":"tool","name":"delegate_task","scope_category":"start","metadata":{"session_id":"docs-parent-session","tool_call_id":"call_delegate"},"data":{"goal":"Run the command `printf docs_nested_leaf_function` using the terminal tool.","toolsets":["terminal"]}} -{"kind":"mark","name":"hermes.subagent.start","metadata":{"parent_session_id":"docs-parent-session","session_id":"docs-child-session","subagent_id":"sa-0-docs","child_role":"leaf"}} -{"kind":"scope","category":"tool","name":"terminal","scope_category":"end","metadata":{"session_id":"docs-child-session","tool_call_id":"call_terminal","status":"ok"},"data":"{\"output\":\"docs_nested_leaf_function\",\"exit_code\":0,\"error\":null}"} -{"kind":"scope","category":"tool","name":"delegate_task","scope_category":"end","metadata":{"session_id":"docs-parent-session","tool_call_id":"call_delegate","status":"ok"}} -``` - -Sanitized ATIF excerpt: - -```json -{ - "schema_version": "ATIF-v1.7", - "session_id": "docs-parent-session", - "agent": {"name": "Hermes Agent E2E", "version": "docs-example", "model_name": "qwen3.6:35b"}, - "steps": [ - { - "source": "agent", - "tool_calls": [{"function_name": "delegate_task"}], - "observation": { - "results": [ - { - "subagent_trajectory_ref": [{"session_id": "docs-child-session"}], - "content": "{\"results\":[{\"status\":\"completed\",\"tool_trace\":[{\"tool\":\"terminal\",\"status\":\"ok\"}]}]}" - } - ] - } - }, - {"source": "agent", "message": "parent received nested subagent result."} - ], - "subagent_trajectories": [ - { - "session_id": "docs-child-session", - "steps": [ - { - "source": "agent", - "tool_calls": [{"function_name": "terminal", "arguments": {"command": "printf docs_nested_leaf_function"}}], - "observation": {"results": [{"content": "{\"output\":\"docs_nested_leaf_function\",\"exit_code\":0,\"error\":null}"}]} - } - ] - } - ] -} -``` - -### Parallel Tool Calls - -This run asks the model to emit two `read_file` tool calls in the same assistant -message. Hermes dispatches the read-only tools as one batch, and NeMo Relay -records both tool invocations. - -```bash -mkdir -p /tmp/hermes-nemo-relay-docs/workdir -printf 'docs_parallel_alpha_function\n' > /tmp/hermes-nemo-relay-docs/workdir/alpha.txt -printf 'docs_parallel_beta_function\n' > /tmp/hermes-nemo-relay-docs/workdir/beta.txt -cd /tmp/hermes-nemo-relay-docs/workdir - -export HERMES_NEMO_RELAY_ATOF_ENABLED=1 -export HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY=/tmp/hermes-nemo-relay-docs/parallel/atof -export HERMES_NEMO_RELAY_ATOF_FILENAME=parallel-tools-atof.jsonl -export HERMES_NEMO_RELAY_ATOF_MODE=overwrite -export HERMES_NEMO_RELAY_ATIF_ENABLED=1 -export HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY=/tmp/hermes-nemo-relay-docs/parallel/atif -export HERMES_NEMO_RELAY_ATIF_FILENAME_TEMPLATE='parallel-tools-atif-{session_id}.json' -export HERMES_NEMO_RELAY_ATIF_AGENT_NAME='Hermes Agent E2E' -export HERMES_NEMO_RELAY_ATIF_AGENT_VERSION=docs-example - -hermes chat \ - --query 'Use exactly two read_file tool calls in the same assistant message. Read alpha.txt and beta.txt. Do not call terminal. After both tool results are available, reply with exactly: parallel tools complete.' \ - --provider custom \ - --model qwen3.6:35b \ - --toolsets file \ - --max-turns 8 \ - --quiet \ - --accept-hooks -``` - -CLI output: - -```text -session_id: docs-parallel-session -parallel tools complete. -``` - -Sanitized ATOF excerpt: - -```jsonl -{"kind":"scope","category":"llm","name":"custom","scope_category":"end","data":{"assistant_message":{"tool_calls":[{"id":"call_alpha","name":"read_file","arguments":"{\"path\":\"alpha.txt\"}"},{"id":"call_beta","name":"read_file","arguments":"{\"path\":\"beta.txt\"}"}]},"finish_reason":"tool_calls"}} -{"kind":"scope","category":"tool","name":"read_file","scope_category":"start","timestamp":"2026-05-31T00:15:08.956732+00:00","metadata":{"session_id":"docs-parallel-session","tool_call_id":"call_alpha"},"data":{"path":"alpha.txt"}} -{"kind":"scope","category":"tool","name":"read_file","scope_category":"start","timestamp":"2026-05-31T00:15:08.956804+00:00","metadata":{"session_id":"docs-parallel-session","tool_call_id":"call_beta"},"data":{"path":"beta.txt"}} -{"kind":"scope","category":"tool","name":"read_file","scope_category":"end","metadata":{"session_id":"docs-parallel-session","tool_call_id":"call_beta","status":"ok"},"data":"{\"content\":\" 1|docs_parallel_beta_function\\n\"}"} -{"kind":"scope","category":"tool","name":"read_file","scope_category":"end","metadata":{"session_id":"docs-parallel-session","tool_call_id":"call_alpha","status":"ok"},"data":"{\"content\":\" 1|docs_parallel_alpha_function\\n\"}"} -``` - -Sanitized ATIF excerpt: - -```json -{ - "schema_version": "ATIF-v1.7", - "session_id": "docs-parallel-session", - "agent": {"name": "Hermes Agent E2E", "version": "docs-example", "model_name": "qwen3.6:35b"}, - "steps": [ - { - "source": "agent", - "tool_calls": [ - {"tool_call_id": "call_alpha", "function_name": "read_file", "arguments": {"path": "alpha.txt"}}, - {"tool_call_id": "call_beta", "function_name": "read_file", "arguments": {"path": "beta.txt"}} - ], - "observation": { - "results": [ - {"source_call_id": "call_beta", "content": "{\"content\":\" 1|docs_parallel_beta_function\\n\"}"}, - {"source_call_id": "call_alpha", "content": "{\"content\":\" 1|docs_parallel_alpha_function\\n\"}"} - ] - } - }, - {"source": "agent", "message": "parallel tools complete."} - ] -} -``` - -## ATOF Mapping - -The plugin keeps NeMo Relay's native event model: - -- Hermes sessions map to `agent` scopes. -- Hermes core managed provider calls map to `llm` scope start/end events. -- Hermes core managed tool calls map to `tool` scope start/end events. -- Turn, approval, subagent, and diagnostic fallback events map to `mark` - events. - -For subagent correlation, mark metadata includes parent and child session IDs, -subagent IDs, role/status fields when present, and derived -`parent_trajectory_id` / `child_trajectory_id` values. This keeps the ATOF -stream lossless for later ATIF conversion that can compact subagents into -separate trajectories. - -## Adaptive Execution Example - -Hermes core owns the LLM and tool boundaries and enters NeMo Relay managed -execution while a Hermes-managed Relay consumer is active. With no shared -metrics subscriber or explicitly configured Relay plugin, Hermes calls the -provider or tool directly. The `observability/nemo_relay` plugin retains the -managed path while its adaptive components are installed on those boundaries. - -Minimal `plugins.toml`: - -```toml -version = 1 - -[[components]] -kind = "adaptive" -enabled = true - -[components.config.tool_parallelism] -mode = "observe_only" -``` - -Enable it for Hermes: - -```bash -export HERMES_NEMO_RELAY_PLUGINS_TOML=/tmp/hermes-middleware-test/plugins.toml -``` - -Execution follows these boundaries with or without an adaptive component: - -```text -Hermes provider call - -> nemo_relay.llm.execute(...) - -> Hermes provider adapter callback(...) - -Hermes tool call - -> nemo_relay.tools.execute(...) - -> Hermes authorization and dispatch callback(...) -``` - -The plugin emits observer marks for sessions, turns, approvals, and subagents. -It does not register provider or tool lifecycle hooks, so each managed call -produces one Relay lifecycle. - -### Local Adaptive E2E - -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]`, as provided by NeMo Relay 0.6 -and later. - -```bash -export HERMES_HOME=/tmp/hermes-middleware-test/hermes-home -mkdir -p "$HERMES_HOME" /tmp/hermes-middleware-test/nemo-relay - -cat > "$HERMES_HOME/config.yaml" <<'YAML' -model: - provider: custom - default: qwen3.6:35b - base_url: http://127.0.0.1:11434/v1 - api_key: ollama -plugins: - enabled: - - observability/nemo_relay -YAML - -cat > /tmp/hermes-middleware-test/nemo-relay/plugins.toml <<'TOML' -version = 1 - -[[components]] -kind = "observability" -enabled = true - -[components.config] -version = 1 - -[components.config.atof] -enabled = true -output_directory = "/tmp/hermes-middleware-test/atof" -filename = "middleware-events.jsonl" -mode = "overwrite" - -[components.config.atif] -enabled = true -output_directory = "/tmp/hermes-middleware-test/atif" -filename_template = "middleware-trajectory-{session_id}.json" -agent_name = "Hermes Middleware E2E" -agent_version = "local" - -[[components]] -kind = "adaptive" -enabled = true - -[components.config.tool_parallelism] -mode = "observe_only" -TOML - -export HERMES_NEMO_RELAY_PLUGINS_TOML=/tmp/hermes-middleware-test/nemo-relay/plugins.toml - -hermes chat \ - --query 'Use the terminal tool exactly once to run printf middleware_execution_ok. Then reply with exactly the command output.' \ - --provider custom \ - --model qwen3.6:35b \ - --toolsets terminal \ - --max-turns 4 \ - --quiet \ - --accept-hooks -``` - -Expected CLI output: - -```text -session_id: middleware-demo-session -middleware_execution_ok -``` - -Expected ATOF shape: - -```jsonl -{"kind":"scope","category":"llm","name":"custom","scope_category":"start","metadata":{"session_id":"middleware-demo-session"},"data":{"mode":"observe_only"}} -{"kind":"scope","category":"tool","name":"terminal","scope_category":"start","metadata":{"session_id":"middleware-demo-session","tool_call_id":"call_terminal"},"data":{"mode":"observe_only"}} -{"kind":"scope","category":"tool","name":"terminal","scope_category":"end","metadata":{"session_id":"middleware-demo-session","tool_call_id":"call_terminal","status":"ok"},"data":"{\"output\":\"middleware_execution_ok\",\"exit_code\":0,\"error\":null}"} -``` - -Expected ATIF shape: - -```json -{ - "schema_version": "ATIF-v1.7", - "session_id": "middleware-demo-session", - "agent": { - "name": "Hermes Middleware E2E", - "version": "local", - "model_name": "qwen3.6:35b" - }, - "steps": [ - { - "source": "agent", - "tool_calls": [ - { - "function_name": "terminal", - "arguments": {"command": "printf middleware_execution_ok"} - } - ], - "observation": { - "results": [ - { - "source_call_id": "call_terminal", - "content": "{\"output\":\"middleware_execution_ok\",\"exit_code\":0,\"error\":null}" - } - ] - } - }, - { - "source": "agent", - "message": "middleware_execution_ok" - } - ] -} -``` diff --git a/plugins/observability/nemo_relay/__init__.py b/plugins/observability/nemo_relay/__init__.py deleted file mode 100644 index 13218ed5e4..0000000000 --- a/plugins/observability/nemo_relay/__init__.py +++ /dev/null @@ -1,1023 +0,0 @@ -"""nemo_relay — optional Hermes plugin for NeMo Relay observability.""" - -from __future__ import annotations - -import atexit -import asyncio -import inspect -import json -import logging -import os -import threading -import tomllib -from collections.abc import Callable -from dataclasses import dataclass, field -from pathlib import Path -from typing import Any, Optional - -from agent import relay_runtime - -logger = logging.getLogger(__name__) - -_INIT_FAILED = object() -_LOCK = threading.RLock() -_RUNTIMES: dict[str, "_Runtime | object"] = {} -_SESSION_INITIALIZER_NAME = "hermes.nemo_relay.rich_observability" - - -@dataclass -class _SessionState: - session_id: str - relay_session: relay_runtime.RelaySession | None = None - handle: Any = None - atif_exporter: Any = None - atif_subscriber_name: str = "" - is_embedded_subagent: bool = False - parent_session_id: str = "" - - -@dataclass -class _SubagentContext: - parent_session_id: str - metadata: dict[str, Any] - - -@dataclass -class _Settings: - plugins_toml_path: str = "" - plugins_config: dict[str, Any] | None = None - dynamic_plugins: list[dict[str, Any]] = field(default_factory=list) - atof_enabled: bool = False - atof_output_directory: str = "" - atof_filename: str = "hermes-atof.jsonl" - atof_mode: str = "append" - atif_enabled: bool = False - atif_output_directory: str = "" - atif_filename_template: str = "hermes-atif-{session_id}.json" - atif_subagent_export_mode: str = "embedded" - atif_agent_name: str = "Hermes Agent" - atif_agent_version: str = "unknown" - atif_model_name: str = "unknown" - - -class _ProcessPluginConfiguration: - """Own Relay's process-global plugin configuration across profile runtimes.""" - - def __init__(self) -> None: - self._lock = threading.RLock() - self._key: str | None = None - self._plugin_mod: Any = None - self._activation: Any = None - self._owners: set[int] = set() - - def acquire( - self, - owner: "_Runtime", - plugin_mod: Any, - plugin_config: dict[str, Any], - dynamic_plugins: list[dict[str, Any]], - ) -> tuple[bool, Any]: - owner_id = id(owner) - key = _plugin_configuration_key(plugin_config, dynamic_plugins) - with self._lock: - if owner_id in self._owners: - return True, self._activation - if self._owners: - if self._plugin_mod is plugin_mod and self._key == key: - self._owners.add(owner_id) - return True, self._activation - logger.warning( - "NeMo Relay plugin configuration is already active for another " - "Hermes profile; keeping the existing process-global configuration " - "and using direct observability for this profile." - ) - return False, None - - activation = None - if dynamic_plugins: - initialize_dynamic = getattr( - plugin_mod, - "initialize_with_dynamic_plugins", - None, - ) - if callable(initialize_dynamic): - try: - activation = _resolve_awaitable( - initialize_dynamic(plugin_config, dynamic_plugins) - ) - 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 a binding that exposes " - "plugin.initialize_with_dynamic_plugins (available in NeMo " - "Relay 0.6+). Continuing with static observability only." - ) - - if activation is None: - initialize = getattr(plugin_mod, "initialize", None) - if not callable(initialize): - return False, None - try: - _resolve_awaitable(initialize(plugin_config)) - except Exception as exc: - logger.debug( - "NeMo Relay plugins.toml init failed: %s", - exc, - exc_info=True, - ) - return False, None - - self._key = key - self._plugin_mod = plugin_mod - self._activation = activation - self._owners.add(owner_id) - return True, activation - - def release(self, owner: "_Runtime", nemo_relay: Any) -> None: - owner_id = id(owner) - with self._lock: - if owner_id not in self._owners: - return - self._owners.remove(owner_id) - if self._owners: - return - - failures: list[str] = [] - activation = self._activation - plugin_mod = self._plugin_mod - try: - if activation is not None: - try: - _flush_relay_subscribers(nemo_relay) - except Exception as exc: - failures.append(f"subscriber flush failed: {exc}") - close = getattr(activation, "close", None) - if callable(close): - try: - _resolve_awaitable(close()) - except Exception as exc: - failures.append( - f"dynamic plugin activation close failed: {exc}" - ) - else: - failures.append("dynamic plugin activation has no close method") - else: - clear = getattr(plugin_mod, "clear", None) - if callable(clear): - try: - _resolve_awaitable(clear()) - except Exception as exc: - failures.append( - f"static plugin configuration clear failed: {exc}" - ) - finally: - self._key = None - self._plugin_mod = None - self._activation = None - - if failures: - raise RuntimeError("; ".join(failures)) - - def reset_for_tests(self) -> None: - with self._lock: - self._key = None - self._plugin_mod = None - self._activation = None - self._owners.clear() - - -_PLUGIN_CONFIGURATION = _ProcessPluginConfiguration() - - -class _Runtime: - def __init__( - self, - nemo_relay: Any, - settings: _Settings, - host: relay_runtime.RelayRuntime, - ) -> None: - self.nemo_relay = nemo_relay - self.settings = settings - self.host = host - self._sessions_lock = threading.RLock() - self.sessions: dict[str, _SessionState] = {} - self.subagent_contexts: dict[str, _SubagentContext] = {} - self.atof_exporter: Any = None - self._atof_subscriber_name = f"hermes.nemo_relay.atof.{self.host.runtime_id}" - self._execution_consumer_name = ( - f"hermes.nemo_relay.rich_observability.{self.host.runtime_id}" - ) - self._execution_consumer_retained = False - 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: - self._activate_direct_fallbacks() - self._sync_managed_execution() - - def _sync_managed_execution(self) -> None: - required = bool( - self._plugin_config_initialized - or self.atof_exporter is not None - or self.settings.atif_enabled - ) - if required and not self._execution_consumer_retained: - self.host.retain_managed_execution(self._execution_consumer_name) - self._execution_consumer_retained = True - elif not required and self._execution_consumer_retained: - self.host.release_managed_execution(self._execution_consumer_name) - self._execution_consumer_retained = False - - def _configure_plugins_toml(self) -> bool: - 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) - self._ensure_plugin_config_output_dirs(plugin_config) - initialized, activation = _PLUGIN_CONFIGURATION.acquire( - self, - plugin_mod, - plugin_config, - self.settings.dynamic_plugins, - ) - self._plugin_activation = activation - if activation is not None: - self._ensure_shutdown_registered() - return initialized - - 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 - try: - _PLUGIN_CONFIGURATION.release(self, self.nemo_relay) - finally: - self._plugin_activation = None - self._plugin_config_initialized = False - self._plugin_config_needs_reinit = bool(self.settings.plugins_config) - - def _activate_direct_fallbacks(self) -> None: - self._plugin_config_needs_reinit = False - self._configure_atof() - - def _maybe_reinitialize_plugins_toml(self) -> None: - if not self._plugin_config_needs_reinit or self._plugin_config_initialized: - return - self._plugin_config_initialized = self._configure_plugins_toml() - if not self._plugin_config_initialized: - self._activate_direct_fallbacks() - self._sync_managed_execution() - return - self._clear_atof() - self._plugin_config_needs_reinit = False - self._sync_managed_execution() - - def _plugins_toml_owns_exporter(self, exporter_name: str) -> bool: - return self._plugin_config_initialized and _observability_exporter_enabled( - self.settings.plugins_config, - exporter_name, - ) - - def _ensure_plugin_config_output_dirs(self, config: dict[str, Any]) -> None: - for component in config.get("components", []): - if not isinstance(component, dict): - continue - if component.get("kind") != "observability": - continue - if component.get("enabled") is False: - continue - component_config = component.get("config") - if not isinstance(component_config, dict): - continue - for exporter_name in ("atof", "atif"): - exporter_config = component_config.get(exporter_name) - if not isinstance(exporter_config, dict): - continue - output_directory = exporter_config.get("output_directory") - if isinstance(output_directory, str) and output_directory.strip(): - Path(output_directory).mkdir(parents=True, exist_ok=True) - - def _configure_atof(self) -> None: - if not self.settings.atof_enabled or self.atof_exporter is not None: - return - config = self.nemo_relay.AtofExporterConfig() - if self.settings.atof_output_directory: - Path(self.settings.atof_output_directory).mkdir(parents=True, exist_ok=True) - config.output_directory = self.settings.atof_output_directory - config.filename = self.settings.atof_filename - if self.settings.atof_mode.lower() == "overwrite": - config.mode = self.nemo_relay.AtofExporterMode.Overwrite - else: - config.mode = self.nemo_relay.AtofExporterMode.Append - self.atof_exporter = self.nemo_relay.AtofExporter(config) - self.atof_exporter.register(self._atof_subscriber_name) - - def _clear_atof(self) -> None: - if self.atof_exporter is None: - return - deregister = getattr(self.atof_exporter, "deregister", None) - if callable(deregister): - try: - deregister(self._atof_subscriber_name) - except Exception: - logger.debug("NeMo Relay ATOF deregister failed", exc_info=True) - self.atof_exporter = None - self._sync_managed_execution() - - def prepare_session(self, kwargs: dict[str, Any]) -> _SessionState: - """Register per-session subscribers without opening the core scope.""" - session_id = _session_id(kwargs) - with self._sessions_lock: - self._maybe_reinitialize_plugins_toml() - state = self.sessions.get(session_id) - if state is not None: - return state - - state = _SessionState(session_id=session_id) - if self.settings.atif_enabled and not self._plugins_toml_owns_exporter("atif"): - state.atif_exporter = self.nemo_relay.AtifExporter( - session_id, - self.settings.atif_agent_name, - self.settings.atif_agent_version, - model_name=str(kwargs.get("model") or self.settings.atif_model_name), - extra={ - "source": "hermes-agent", - "plugin": "observability/nemo_relay", - }, - ) - state.atif_subscriber_name = ( - f"hermes.nemo_relay.atif.{self.host.runtime_id}.{session_id}" - ) - state.atif_exporter.register(state.atif_subscriber_name) - self.sessions[session_id] = state - return state - - def ensure_session(self, kwargs: dict[str, Any]) -> _SessionState: - state = self.prepare_session(kwargs) - if state.relay_session is not None: - return state - - rich_metadata = _metadata(kwargs) - with self._sessions_lock: - subagent_context = self.subagent_contexts.get(state.session_id) - if subagent_context is not None: - rich_metadata = {**rich_metadata, **subagent_context.metadata} - relay_session = self.host.ensure_session( - kwargs, - data={"session_id": state.session_id}, - metadata=rich_metadata, - ) - if relay_session is None: - raise RuntimeError("Hermes core Relay session is unavailable") - state.relay_session = relay_session - state.handle = relay_session.handle - if subagent_context is not None: - state.is_embedded_subagent = True - state.parent_session_id = subagent_context.parent_session_id - return state - - def run_in_session( - self, - state: _SessionState, - callback: Callable[..., Any], - *args: Any, - **kwargs: Any, - ) -> Any: - if state.relay_session is None: - raise RuntimeError("Hermes core Relay session is unavailable") - return self.host.run_in_session( - state.relay_session, - callback, - *args, - **kwargs, - ) - - def export_atif(self, state: _SessionState) -> None: - if not self.settings.atif_enabled or state.atif_exporter is None: - return - if state.is_embedded_subagent and self.settings.atif_subagent_export_mode != "all": - return - output_dir = self.settings.atif_output_directory - if not output_dir: - return - Path(output_dir).mkdir(parents=True, exist_ok=True) - filename = self.settings.atif_filename_template.format(session_id=state.session_id) - Path(output_dir, filename).write_text(state.atif_exporter.export_json(), encoding="utf-8") - - def close_session( - self, - kwargs: dict[str, Any], - *, - close_host: bool = True, - ) -> None: - session_id = _session_id(kwargs) - with self._sessions_lock: - self.subagent_contexts.pop(session_id, None) - state = self.sessions.pop(session_id, None) - if state is None: - return - failures: list[str] = [] - if close_host: - try: - self.host.close_session(kwargs) - except Exception as exc: - failures.append(f"core session close failed: {exc}") - try: - self.export_atif(state) - except Exception as exc: - failures.append(f"ATIF export failed: {exc}") - if state.atif_exporter is not None and state.atif_subscriber_name: - try: - state.atif_exporter.deregister(state.atif_subscriber_name) - except Exception as exc: - failures.append(f"ATIF deregister failed: {exc}") - with self._sessions_lock: - if ( - self._plugin_config_initialized - and self._plugin_activation is None - and not self.sessions - ): - try: - self._clear_plugins_toml() - except Exception as exc: - failures.append(f"plugin configuration clear failed: {exc}") - elif ( - self.settings.plugins_config - and self._plugin_activation is None - and not self.sessions - ): - self._plugin_config_needs_reinit = True - if failures: - logger.warning( - "NeMo Relay session %s teardown completed with errors: %s", - session_id, - "; ".join(failures), - ) - - def shutdown(self) -> None: - """Close active sessions and the process-lifetime plugin activation.""" - failures: list[str] = [] - with self._sessions_lock: - session_ids = list(self.sessions) - for session_id in session_ids: - try: - self.close_session({"session_id": session_id, "reason": "runtime_shutdown"}) - except Exception as exc: - failures.append(f"session {session_id} close failed: {exc}") - if self._plugin_config_initialized: - try: - self._clear_plugins_toml() - except Exception as exc: - failures.append(f"plugin runtime close failed: {exc}") - self._clear_atof() - if self._execution_consumer_retained: - self.host.release_managed_execution(self._execution_consumer_name) - self._execution_consumer_retained = False - if self._shutdown_registered and self._plugin_activation is None: - atexit.unregister(self.shutdown) - self._shutdown_registered = False - if failures: - logger.warning( - "NeMo Relay runtime shutdown completed with errors: %s", - "; ".join(failures), - ) - - def mark(self, name: str, kwargs: dict[str, Any]) -> None: - state = self.ensure_session(kwargs) - self.run_in_session( - state, - self.nemo_relay.scope.event, - name, - handle=state.handle, - data=_jsonable(kwargs), - metadata=_metadata(kwargs), - ) - - def mark_subagent_start(self, kwargs: dict[str, Any]) -> None: - parent_state = self.ensure_session(kwargs) - metadata = _metadata(kwargs) - child_session_id = _child_session_id(kwargs) - if child_session_id: - with self._sessions_lock: - self.subagent_contexts[child_session_id] = _SubagentContext( - parent_session_id=parent_state.session_id, - metadata=_subagent_child_metadata(kwargs, metadata), - ) - self.run_in_session( - parent_state, - self.nemo_relay.scope.event, - "hermes.subagent.start", - handle=parent_state.handle, - data=_jsonable(kwargs), - metadata=metadata, - ) - - def mark_subagent_stop(self, kwargs: dict[str, Any]) -> None: - child_session_id = _child_session_id(kwargs) - if child_session_id: - self.close_session( - {"session_id": child_session_id}, - close_host=False, - ) - with self._sessions_lock: - self.subagent_contexts.pop(child_session_id, None) - self.mark("hermes.subagent.stop", kwargs) - -def register(ctx) -> None: - relay_runtime.SESSION_COORDINATOR.register_session_initializer( - _SESSION_INITIALIZER_NAME, - _prepare_core_session, - ) - # 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) - ctx.register_hook("on_session_reset", on_session_reset) - ctx.register_hook("pre_llm_call", on_pre_llm_call) - ctx.register_hook("post_llm_call", on_post_llm_call) - ctx.register_hook("pre_approval_request", on_pre_approval_request) - ctx.register_hook("post_approval_response", on_post_approval_response) - ctx.register_hook("subagent_start", on_subagent_start) - ctx.register_hook("subagent_stop", on_subagent_stop) - - -def on_session_start(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.ensure_session(kwargs)) - - -def on_session_end(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: (runtime.mark("hermes.session.end", kwargs), runtime.export_atif(runtime.ensure_session(kwargs)))) - - -def on_session_finalize(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.close_session(kwargs, close_host=False)) - - -def on_session_reset(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.close_session(kwargs, close_host=False)) - - -def on_pre_llm_call(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.mark("hermes.turn.start", kwargs)) - - -def on_post_llm_call(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.mark("hermes.turn.end", kwargs)) - - -def on_pre_approval_request(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.mark("hermes.approval.request", kwargs)) - - -def on_post_approval_response(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.mark("hermes.approval.response", kwargs)) - - -def on_subagent_start(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.mark_subagent_start(kwargs)) - - -def on_subagent_stop(**kwargs: Any) -> None: - runtime = _get_runtime() - if runtime is not None: - _safe(lambda: runtime.mark_subagent_stop(kwargs)) - - -def _prepare_core_session( - host: relay_runtime.RelayRuntime, - context: dict[str, Any], -) -> None: - """Register rich subscribers before core creates the conversation scope.""" - runtime = _get_runtime( - profile_key=str(context.get("profile_key") or host.profile_key), - host=host, - ) - if runtime is not None: - runtime.prepare_session(context) - - -def _get_runtime( - *, - profile_key: str | None = None, - host: relay_runtime.RelayRuntime | None = None, -) -> Optional[_Runtime]: - profile_key = profile_key or relay_runtime.current_profile_key() - with _LOCK: - runtime = _RUNTIMES.get(profile_key) - if runtime is _INIT_FAILED: - return None - if isinstance(runtime, _Runtime): - if host is None or runtime.host is host: - return runtime - runtime.shutdown() - _RUNTIMES.pop(profile_key, None) - try: - resolved_host = host or relay_runtime.get_runtime(profile_key=profile_key) - if resolved_host is None: - raise RuntimeError("Hermes core Relay runtime is unavailable") - runtime = _Runtime( - nemo_relay=resolved_host.relay, - settings=_load_settings(), - host=resolved_host, - ) - except Exception as exc: - logger.debug("NeMo Relay plugin disabled: init failed: %s", exc, exc_info=True) - _RUNTIMES[profile_key] = _INIT_FAILED - return None - _RUNTIMES[profile_key] = runtime - return runtime - - -def _load_settings() -> _Settings: - plugins_toml_path = _env("HERMES_NEMO_RELAY_PLUGINS_TOML") - plugins_config = _load_plugins_config(plugins_toml_path) - return _Settings( - plugins_toml_path=plugins_toml_path, - plugins_config=plugins_config, - dynamic_plugins=_dynamic_plugin_specs(plugins_config, plugins_toml_path), - atof_enabled=_env_bool("HERMES_NEMO_RELAY_ATOF_ENABLED"), - atof_output_directory=_env("HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY"), - atof_filename=_env("HERMES_NEMO_RELAY_ATOF_FILENAME") or "hermes-atof.jsonl", - atof_mode=_env("HERMES_NEMO_RELAY_ATOF_MODE") or "append", - atif_enabled=_env_bool("HERMES_NEMO_RELAY_ATIF_ENABLED"), - atif_output_directory=_env("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY"), - atif_filename_template=_env("HERMES_NEMO_RELAY_ATIF_FILENAME_TEMPLATE") or "hermes-atif-{session_id}.json", - atif_subagent_export_mode=_atif_subagent_export_mode(), - atif_agent_name=_env("HERMES_NEMO_RELAY_ATIF_AGENT_NAME") or "Hermes Agent", - atif_agent_version=_env("HERMES_NEMO_RELAY_ATIF_AGENT_VERSION") or "unknown", - atif_model_name=_env("HERMES_NEMO_RELAY_ATIF_MODEL_NAME") or "unknown", - ) - - -def _static_plugin_config(plugins_config: dict[str, Any]) -> dict[str, Any]: - """Return Relay's base config without embedding- or gateway-host fields.""" - return { - key: value - for key, value in plugins_config.items() - if key not in {"dynamic_plugins", "plugins"} - } - - -def _plugin_configuration_key( - plugin_config: dict[str, Any], - dynamic_plugins: list[dict[str, Any]], -) -> str: - return json.dumps( - {"config": plugin_config, "dynamic_plugins": dynamic_plugins}, - sort_keys=True, - separators=(",", ":"), - default=str, - ) - - -def _dynamic_plugin_specs( - plugins_config: dict[str, Any] | None, - plugins_toml_path: str = "", -) -> list[dict[str, Any]]: - if not isinstance(plugins_config, dict): - return [] - - raw_specs = plugins_config.get("dynamic_plugins") - plugins_section = plugins_config.get("plugins") - if plugins_section is not None: - if not isinstance(plugins_section, dict): - logger.error( - "Invalid NeMo Relay plugins config: expected [plugins] to be an object; " - "no dynamic plugins will be activated. Continuing with static " - "observability only." - ) - return [] - if plugins_section: - logger.error( - "Hermes cannot activate Relay gateway [[plugins.dynamic]] records because " - "the Python binding does not expose the CLI lifecycle resolver for " - "enablement, trust policy, and worker environments. Use Hermes-owned " - "[[dynamic_plugins]] activation specs instead; no dynamic plugins will be " - "activated. Continuing with static observability only." - ) - return [] - 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]] = [] - invalid = False - for index, raw_spec in enumerate(raw_specs): - if not isinstance(raw_spec, dict): - logger.warning( - "Invalid NeMo Relay dynamic_plugins[%d]: expected an object", index - ) - invalid = True - 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( - "Invalid NeMo Relay dynamic_plugins[%d]: plugin_id is required", index - ) - invalid = True - continue - if kind not in {"rust_dynamic", "worker"}: - logger.warning( - "Invalid NeMo Relay dynamic_plugins[%d]: kind must be rust_dynamic or worker", - index, - ) - invalid = True - continue - if not isinstance(manifest_ref, str) or not manifest_ref.strip(): - logger.warning( - "Invalid NeMo Relay dynamic_plugins[%d]: manifest_ref is required", index - ) - invalid = True - continue - if not isinstance(config, dict): - logger.warning( - "Invalid NeMo Relay dynamic_plugins[%d]: config must be an object", index - ) - invalid = True - continue - if environment_ref is not None and ( - not isinstance(environment_ref, str) or not environment_ref.strip() - ): - logger.warning( - "Invalid NeMo Relay dynamic_plugins[%d]: environment_ref must be a " - "non-empty string", - index, - ) - invalid = True - continue - spec: dict[str, Any] = { - "plugin_id": plugin_id.strip(), - "kind": kind, - "manifest_ref": _config_relative_path(manifest_ref.strip(), plugins_toml_path), - "config": config, - } - if environment_ref is not None: - spec["environment_ref"] = _config_relative_path( - environment_ref.strip(), plugins_toml_path - ) - specs.append(spec) - if invalid: - logger.error( - "NeMo Relay dynamic plugin configuration is invalid; no dynamic plugins " - "will be activated. Continuing with static observability only." - ) - return [] - return specs - - -def _config_relative_path(value: str, plugins_toml_path: str) -> str: - """Resolve a plugin path relative to its physical ``plugins.toml`` file.""" - path = Path(value) - if path.is_absolute(): - return str(path) - config_path = Path(plugins_toml_path) if plugins_toml_path else Path.cwd() / "plugins.toml" - if not config_path.is_absolute(): - config_path = Path.cwd() / config_path - return os.path.abspath(config_path.parent / path) - - -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 - try: - return tomllib.loads(Path(path).read_text(encoding="utf-8")) - except Exception as exc: - logger.debug("NeMo Relay plugins.toml load failed: %s", exc, exc_info=True) - return None - - -def _enabled_component_config( - plugins_config: dict[str, Any] | None, - kind: str, -) -> dict[str, Any] | None: - if not isinstance(plugins_config, dict): - return None - components = plugins_config.get("components") - if not isinstance(components, list): - return None - for component in components: - if not isinstance(component, dict): - continue - if component.get("kind") != kind or not component.get("enabled", True): - continue - config = component.get("config") - return config if isinstance(config, dict) else {} - return None - - -def _observability_exporter_enabled( - plugins_config: dict[str, Any] | None, - exporter_name: str, -) -> bool: - observability_config = _enabled_component_config(plugins_config, "observability") - if not isinstance(observability_config, dict): - return False - exporter_config = observability_config.get(exporter_name) - if not isinstance(exporter_config, dict): - return False - return exporter_config.get("enabled", True) is not False - - -def _env(name: str) -> str: - return os.environ.get(name, "").strip() - - -def _atif_subagent_export_mode() -> str: - mode = _env("HERMES_NEMO_RELAY_ATIF_SUBAGENT_EXPORT_MODE").lower() - return "all" if mode == "all" else "embedded" - - -def _env_bool(name: str) -> bool: - return _env(name).lower() in {"1", "true", "yes", "on"} - - -def _session_id(kwargs: dict[str, Any]) -> str: - return str(kwargs.get("session_id") or kwargs.get("parent_session_id") or "default") - - -def _child_session_id(kwargs: dict[str, Any]) -> str: - return str(kwargs.get("child_session_id") or "") - - -def _subagent_child_metadata( - kwargs: dict[str, Any], - parent_metadata: dict[str, Any], -) -> dict[str, Any]: - child_session_id = _child_session_id(kwargs) - metadata = { - "session_id": child_session_id, - "trajectory_id": child_session_id, - "nemo_relay_scope_role": "subagent", - } - for target, source in ( - ("subagent_id", "child_subagent_id"), - ("child_session_id", "child_session_id"), - ("child_subagent_id", "child_subagent_id"), - ("child_role", "child_role"), - ("parent_session_id", "parent_session_id"), - ("parent_turn_id", "parent_turn_id"), - ("parent_subagent_id", "parent_subagent_id"), - ("parent_trajectory_id", "parent_trajectory_id"), - ("telemetry_schema_version", "telemetry_schema_version"), - ): - value = parent_metadata.get(source) - if value is not None: - metadata[target] = value - return metadata - - -def _metadata(kwargs: dict[str, Any]) -> dict[str, Any]: - keys = ( - "telemetry_schema_version", - "session_id", - "platform", - "task_id", - "turn_id", - "api_request_id", - "tool_call_id", - "parent_session_id", - "parent_turn_id", - "parent_subagent_id", - "child_session_id", - "child_subagent_id", - "child_role", - "child_status", - "provider", - "model", - "api_mode", - "status", - "reason", - ) - metadata = { - key: _jsonable(kwargs[key]) - for key in keys - if key in kwargs and kwargs[key] is not None - } - if "session_id" in metadata: - metadata.setdefault("trajectory_id", metadata["session_id"]) - if "parent_session_id" in metadata: - metadata.setdefault("parent_trajectory_id", metadata["parent_session_id"]) - if "child_session_id" in metadata: - metadata.setdefault("child_trajectory_id", metadata["child_session_id"]) - return metadata - - -def _jsonable(value: Any) -> Any: - if value is None or isinstance(value, (str, int, float, bool)): - return value - if isinstance(value, dict): - return {str(k): _jsonable(v) for k, v in value.items()} - if isinstance(value, (list, tuple, set)): - return [_jsonable(v) for v in value] - try: - if hasattr(value, "model_dump"): - return _jsonable(value.model_dump(mode="json")) - except Exception: - pass - try: - if hasattr(value, "__dict__"): - return _jsonable(vars(value)) - except Exception: - pass - try: - return json.loads(json.dumps(value, default=str)) - except Exception: - return str(value) - - -def _safe(fn) -> None: - try: - fn() - except Exception as exc: - logger.debug("NeMo Relay hook handling failed: %s", exc, exc_info=True) - - -def _resolve_awaitable(value: Any) -> Any: - if not inspect.isawaitable(value): - return value - try: - asyncio.get_running_loop() - except RuntimeError: - return asyncio.run(value) - - result: dict[str, Any] = {} - error: dict[str, BaseException] = {} - - def _runner() -> None: - try: - result["value"] = asyncio.run(value) - except BaseException as exc: # pragma: no cover - re-raised below - error["exc"] = exc - - thread = threading.Thread( - target=_runner, - name="hermes-nemo-relay-awaitable", - daemon=True, - ) - thread.start() - thread.join() - if "exc" in error: - raise error["exc"] - return result.get("value") - - -def reset_for_tests() -> None: - relay_runtime.SESSION_COORDINATOR.unregister_session_initializer( - _SESSION_INITIALIZER_NAME - ) - with _LOCK: - runtimes = list(_RUNTIMES.values()) - _RUNTIMES.clear() - for runtime in runtimes: - if isinstance(runtime, _Runtime): - runtime.shutdown() - _PLUGIN_CONFIGURATION.reset_for_tests() diff --git a/plugins/observability/nemo_relay/plugin.yaml b/plugins/observability/nemo_relay/plugin.yaml deleted file mode 100644 index 046d5d0d85..0000000000 --- a/plugins/observability/nemo_relay/plugin.yaml +++ /dev/null @@ -1,15 +0,0 @@ -name: nemo_relay -version: "0.1.0" -description: "Optional NeMo Relay observability for Hermes. Opt in with `hermes plugins enable observability/nemo_relay`; HERMES_NEMO_RELAY_* env vars configure exports after the plugin is enabled." -author: NousResearch -hooks: - - on_session_start - - on_session_end - - on_session_finalize - - on_session_reset - - pre_llm_call - - post_llm_call - - pre_approval_request - - post_approval_response - - subagent_start - - subagent_stop diff --git a/scripts/toolperf_abeval/README.md b/scripts/toolperf_abeval/README.md index 0f8525eaac..5a781d0b32 100644 --- a/scripts/toolperf_abeval/README.md +++ b/scripts/toolperf_abeval/README.md @@ -40,9 +40,11 @@ in real production traffic. provider: openrouter YAML printf 'OPENROUTER_API_KEY=%s\n' "$KEY" > "$ABEVAL_HOME/.env" - HERMES_HOME=$ABEVAL_HOME hermes plugins enable observability/nemo_relay ``` + The runner writes a per-run Relay `plugins.toml` and points the native SDK + integration at it; no Hermes observability plugin needs to be enabled. + 2. Prepare the two trees: ```bash diff --git a/scripts/toolperf_abeval/ab_eval.py b/scripts/toolperf_abeval/ab_eval.py index 176685ecd3..01fa9a355b 100644 --- a/scripts/toolperf_abeval/ab_eval.py +++ b/scripts/toolperf_abeval/ab_eval.py @@ -165,14 +165,34 @@ def run(arm: str, model: str, reps: int, pythonpath: str, only=None): work.mkdir(parents=True, exist_ok=True) make_sandbox(work) atof = resdir / f"{run_id}.atof.jsonl" + relay_config = work / "relay-plugins.toml" + relay_config.write_text( + f""" +version = 1 + +[[components]] +kind = "observability" +enabled = true + +[components.config] +version = 3 + +[components.config.atof] +enabled = true + +[[components.config.atof.sinks]] +type = "file" +output_directory = {json.dumps(str(atof.parent))} +filename = {json.dumps(atof.name)} +mode = "overwrite" +""".strip(), + encoding="utf-8", + ) env = dict(os.environ) env.update({ "PYTHONPATH": pythonpath, "HERMES_HOME": str(HOME), - "HERMES_NEMO_RELAY_ATOF_ENABLED": "1", - "HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY": str(atof.parent), - "HERMES_NEMO_RELAY_ATOF_FILENAME": atof.name, - "HERMES_NEMO_RELAY_ATOF_MODE": "overwrite", + "HERMES_NEMO_RELAY_PLUGINS_TOML": str(relay_config), }) q = TASKS[name].replace("{WORK}", str(work)) t0 = time.time() diff --git a/tests/hermes_cli/test_plugins_cmd_enable_disable_nested.py b/tests/hermes_cli/test_plugins_cmd_enable_disable_nested.py index a964626e23..38e1b57a7c 100644 --- a/tests/hermes_cli/test_plugins_cmd_enable_disable_nested.py +++ b/tests/hermes_cli/test_plugins_cmd_enable_disable_nested.py @@ -3,7 +3,7 @@ Companion to test_plugins_cmd_category_discovery.py. That file covers the *listing* side of nested category plugins (issue #41066). These tests cover the *mutation* side: `hermes plugins enable/disable` must resolve a bare name -OR a full path-derived key (e.g. `observability/nemo_relay`) to the canonical +OR a full path-derived key (e.g. `observability/trace_sink`) to the canonical registry key and write THAT — the same string PluginManager gates on — so a nested bundled plugin can actually be toggled. """ @@ -32,8 +32,8 @@ def _make_category_plugin(parent: Path, category: str, name: str, manifest: dict def nested_plugin_env(tmp_path): """A user-plugins dir containing one nested and one flat plugin, with the bundled dir pointed at an empty path. Returns the tmp_path.""" - _make_category_plugin(tmp_path, "observability", "nemo_relay", { - "name": "nemo_relay", "version": "1.0.0", "description": "relay obs" + _make_category_plugin(tmp_path, "observability", "trace_sink", { + "name": "trace_sink", "version": "1.0.0", "description": "trace sink" }) _make_plugin_dir(tmp_path, "disk-cleanup", { "name": "disk-cleanup", "version": "1.0.0" @@ -53,7 +53,7 @@ class TestResolvePluginKey: from hermes_cli.plugins_cmd import _resolve_plugin_key mock_user.return_value = nested_plugin_env mock_bundled.return_value = nested_plugin_env / "nonexistent" - assert _resolve_plugin_key("observability/nemo_relay") == "observability/nemo_relay" + assert _resolve_plugin_key("observability/trace_sink") == "observability/trace_sink" @patch("hermes_cli.plugins.get_bundled_plugins_dir") @@ -98,13 +98,13 @@ class TestEnableDisableNested: mock_user.return_value = nested_plugin_env mock_bundled.return_value = nested_plugin_env / "nonexistent" - cmd_enable("nemo_relay", allow_tool_override=False) # bare name + cmd_enable("trace_sink", allow_tool_override=False) # bare name saved = mock_save_en.call_args[0][0] # The canonical key — NOT the bare name — must be persisted, because # that is what PluginManager matches when deciding to load. - assert "observability/nemo_relay" in saved - assert "nemo_relay" not in saved or "observability/nemo_relay" in saved + assert "observability/trace_sink" in saved + assert "trace_sink" not in saved or "observability/trace_sink" in saved @patch("hermes_cli.plugins.get_bundled_plugins_dir") diff --git a/tests/plugins/test_nemo_relay_plugin.py b/tests/plugins/test_nemo_relay_plugin.py deleted file mode 100644 index 343bd76965..0000000000 --- a/tests/plugins/test_nemo_relay_plugin.py +++ /dev/null @@ -1,490 +0,0 @@ -"""Tests for the bundled observability/nemo_relay plugin.""" - -from __future__ import annotations - -import asyncio -import contextvars -import gc -import importlib -import json -import sys -import warnings -from pathlib import Path -from types import SimpleNamespace - -import pytest -import yaml - -from hermes_cli import lifecycle, plugins as plugin_api -from hermes_cli.observability import relay_runtime, relay_shared_metrics -from hermes_cli.plugins import PluginManager - - -REPO_ROOT = Path(__file__).resolve().parents[2] -PLUGIN_DIR = REPO_ROOT / "plugins" / "observability" / "nemo_relay" - - -class _FakeNemoRelay: - def __init__(self): - self.events = [] - self._callbacks = {} - self._llm_starts = {} - self._scope_serial = 0 - self._scope_context = contextvars.ContextVar( - "fake_nemo_relay_scope", default=None - ) - self.ScopeType = SimpleNamespace(Agent="agent", Function="function") - self.scope = SimpleNamespace( - push=self._scope_push, - pop=self._scope_pop, - event=self._scope_event, - ) - self.llm = SimpleNamespace( - call=self._llm_call, - call_end=self._llm_call_end, - execute=self._llm_execute, - ) - self.tools = SimpleNamespace( - call=self._tool_call, - call_end=self._tool_call_end, - execute=self._tool_execute, - request_intercepts=self._tool_request_intercepts, - ) - self.plugin = SimpleNamespace( - initialize=self._plugin_initialize, - clear=self._plugin_clear, - initialize_with_dynamic_plugins=self._plugin_initialize_with_dynamic, - ) - self.subscribers = SimpleNamespace( - register=self._register_subscriber, - deregister=self._deregister_subscriber, - flush=self._flush_subscribers, - ) - self.LLMRequest = _FakeLLMRequest - self.AtofExporterConfig = _FakeAtofExporterConfig - self.AtofExporterMode = SimpleNamespace(Append="append", Overwrite="overwrite") - self.AtofExporter = self._make_atof_exporter - self.AtifExporter = self._make_atif_exporter - self.get_scope_stack = self._get_scope_stack - - def _scope_push(self, name, scope_type, **kwargs): - self._scope_serial += 1 - handle = ("scope", name, self._scope_serial) - self._scope_context.set(handle) - self.events.append(("scope.push", name, scope_type, kwargs)) - return handle - - def _scope_pop(self, handle, **kwargs): - self.events.append(("scope.pop", handle, kwargs)) - - def _scope_event(self, name, **kwargs): - self.events.append(("scope.event", name, kwargs)) - - def _get_scope_stack(self): - current = self._scope_context.get() - self.events.append(("scope.sync", current)) - return current - - def _llm_call(self, name, request, **kwargs): - handle = ("llm", name) - self._llm_starts[handle] = kwargs - self.events.append(("llm.call", name, request.content, kwargs)) - return handle - - def _llm_call_end(self, handle, response, **kwargs): - self.events.append(("llm.call_end", handle, response, kwargs)) - start = self._llm_starts.pop(handle, {}) - event = SimpleNamespace( - kind="scope", - category="llm", - name=handle[1], - scope_category="end", - category_profile={"model_name": start.get("model_name")}, - metadata={ - **(start.get("metadata") or {}), - **(kwargs.get("metadata") or {}), - "otel.status_code": "OK", - }, - data=response, - ) - for callback in list(self._callbacks.values()): - callback(event) - - def _llm_execute(self, name, request, func, **kwargs): - self.events.append(("llm.execute.start", name, request.content, kwargs)) - handle = self._llm_call(name, request, **kwargs) - result = func(_FakeLLMRequest(request.headers, {"intercepted": True, **request.content})) - self._llm_call_end( - handle, - result, - **{key: value for key, value in kwargs.items() if key != "handle"}, - ) - self.events.append(("llm.execute.end", name, result, kwargs)) - return result - - def _tool_call(self, name, args, **kwargs): - handle = ("tool", name) - self.events.append(("tool.call", name, args, kwargs)) - return handle - - def _tool_call_end(self, handle, result, **kwargs): - self.events.append(("tool.call_end", handle, result, kwargs)) - - def _tool_execute(self, name, args, func, **kwargs): - self.events.append(("tool.execute.start", name, args, kwargs)) - handle = self._tool_call(name, args, **kwargs) - result = func(args) - self._tool_call_end( - handle, - result, - **{key: value for key, value in kwargs.items() if key != "handle"}, - ) - self.events.append(("tool.execute.end", name, result, kwargs)) - return result - - def _tool_request_intercepts(self, name, args): - self.events.append(("tool.request_intercepts", name, args)) - return {"intercepted": True, **args} - - def _make_atof_exporter(self, config): - return _FakeAtofExporter(self.events, config) - - def _make_atif_exporter(self, session_id, agent_name, agent_version, **kwargs): - return _FakeAtifExporter(self.events, session_id, agent_name, agent_version, kwargs) - - async def _plugin_initialize(self, config): - self.events.append(("plugin.initialize", config)) - return {"diagnostics": []} - - async def _plugin_clear(self): - self.events.append(("plugin.clear",)) - - async def _plugin_initialize_with_dynamic(self, config, dynamic_plugins): - self.events.append(("plugin.activate_dynamic", config, dynamic_plugins)) - return _FakePluginActivation(self.events) - - def _register_subscriber(self, name, callback): - self._callbacks[name] = callback - self.events.append(("subscribers.register", name)) - - def _deregister_subscriber(self, name): - self._callbacks.pop(name, None) - self.events.append(("subscribers.deregister", name)) - - def _flush_subscribers(self): - self.events.append(("subscribers.flush",)) - - -class _FakePluginActivation: - def __init__(self, events): - self.events = events - self.report = {"diagnostics": []} - - async def close(self): - self.events.append(("plugin.activation.close",)) - - -class _FakeLLMRequest: - def __init__(self, headers, content): - self.headers = headers - self.content = content - - -class _FakeAtofExporterConfig: - def __init__(self): - self.output_directory = "" - self.filename = "events.jsonl" - self.mode = "append" - - -class _FakeAtofExporter: - def __init__(self, events, config): - self.events = events - self.config = config - - def register(self, name): - self.events.append(("atof.register", name, self.config.output_directory, self.config.filename)) - - def deregister(self, name): - self.events.append(("atof.deregister", name, self.config.output_directory, self.config.filename)) - return True - - -class _FakeAtifExporter: - def __init__(self, events, session_id, agent_name, agent_version, kwargs): - self.events = events - self.session_id = session_id - self.agent_name = agent_name - self.agent_version = agent_version - self.kwargs = kwargs - - def register(self, name): - self.events.append(("atif.register", name, self.session_id)) - - def deregister(self, name): - self.events.append(("atif.deregister", name, self.session_id)) - return True - - def export_json(self): - self.events.append(("atif.export", self.session_id)) - return json.dumps({"session_id": self.session_id, "agent_name": self.agent_name}) - - -def _fresh_plugin(monkeypatch, fake): - existing = sys.modules.get("plugins.observability.nemo_relay") - if existing is not None: - existing.reset_for_tests() - relay_shared_metrics._reset_for_tests() - relay_runtime._reset_for_tests() - monkeypatch.setattr(relay_runtime, "_load_nemo_relay", lambda: fake) - monkeypatch.setitem(sys.modules, "nemo_relay", fake) - sys.modules.pop("plugins.observability.nemo_relay", None) - plugin = importlib.import_module("plugins.observability.nemo_relay") - plugin.reset_for_tests() - return plugin - - -def _enable_dynamic_plugin(tmp_path, monkeypatch) -> Path: - plugins_toml = tmp_path / "plugins.toml" - plugins_toml.write_text( - f""" -version = 1 - -[[dynamic_plugins]] -plugin_id = "fixture" -kind = "rust_dynamic" -manifest_ref = "{(tmp_path / "fixture" / "relay-plugin.toml").as_posix()}" - -[dynamic_plugins.config] -mode = "test" -""", - encoding="utf-8", - ) - monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml)) - return plugins_toml - - -def test_manifest_fields(): - data = yaml.safe_load((PLUGIN_DIR / "plugin.yaml").read_text()) - assert data["name"] == "nemo_relay" - assert set(data["hooks"]) == { - "on_session_start", - "on_session_end", - "on_session_finalize", - "on_session_reset", - "pre_llm_call", - "post_llm_call", - "pre_approval_request", - "post_approval_response", - "subagent_start", - "subagent_stop", - } - - -def test_nemo_relay_plugin_is_discoverable_as_bundled_plugin(tmp_path, monkeypatch): - monkeypatch.setenv("HERMES_HOME", str(tmp_path / "hermes_test")) - - manager = PluginManager() - manager.discover_and_load() - - loaded = manager._plugins["observability/nemo_relay"] - assert loaded.manifest.name == "nemo_relay" - assert loaded.manifest.source == "bundled" - assert not loaded.enabled - - -def test_shared_metrics_and_rich_plugin_share_one_core_session( - tmp_path, - monkeypatch, -): - from agent import relay_llm - - fake = _FakeNemoRelay() - hermes_home = tmp_path / "hermes-home" - monkeypatch.setenv("HERMES_HOME", str(hermes_home)) - monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif") - ) - monkeypatch.setattr( - "hermes_cli.config.read_raw_config_readonly", - lambda: {"telemetry": {"shared_metrics": {"enabled": True}}}, - ) - plugin = _fresh_plugin(monkeypatch, fake) - manager = PluginManager() - - class _Context: - def register_hook(self, name, callback): - manager._hooks.setdefault(name, []).append(callback) - - plugin.register(_Context()) - monkeypatch.setattr(plugin_api, "_plugin_manager", manager) - - event = { - "session_id": "s1", - "task_id": "t1", - "api_request_id": "api-1", - "provider": "anthropic", - "model": "claude-sonnet", - "platform": "cli", - } - coordinator = relay_runtime.SESSION_COORDINATOR - lease = coordinator.acquire_conversation( - profile_key=relay_runtime.current_profile_key(), - session_id="s1", - platform="cli", - model=event["model"], - ) - lifecycle.invoke_hook("on_session_start", **event) - turn = coordinator.begin_turn( - lease, - turn_id="turn-1", - task_id="t1", - ) - lifecycle.invoke_hook( - "pre_api_request", - **event, - request={"body": {"messages": [{"role": "user", "content": "hi"}]}}, - ) - relay_llm.execute( - {"messages": [{"role": "user", "content": "hi"}]}, - lambda _request: { - "assistant_message": {"role": "assistant", "content": "hello"} - }, - session_id="s1", - name="anthropic", - model_name="claude-sonnet", - metadata={"api_request_id": "api-1", "api_mode": "custom"}, - ) - lifecycle.invoke_hook( - "post_api_request", - **event, - response={"assistant_message": {"role": "assistant", "content": "hello"}}, - ) - coordinator.end_turn(turn, outcome="success") - coordinator.release_conversation(lease) - lifecycle.finalize_session(session_id="s1") - - session_pushes = [ - item - for item in fake.events - if item[0] == "scope.push" and item[1] == relay_runtime.SESSION_SCOPE - ] - assert len(session_pushes) == 1 - register_metrics = next( - index - for index, item in enumerate(fake.events) - if item[0] == "subscribers.register" - and item[1].startswith("hermes.nemo_relay.shared_metrics.") - ) - register_atif = next( - index for index, item in enumerate(fake.events) if item[0] == "atif.register" - ) - open_session = fake.events.index(session_pushes[0]) - assert register_metrics < register_atif < open_session - - plugin_runtime = plugin._get_runtime() - assert plugin_runtime is not None - assert not plugin_runtime.sessions - assert relay_runtime.get_session_handle("s1") is None - packages = list( - (hermes_home / "telemetry" / "shared_metrics" / "outbox").glob("*.json") - ) - assert len(packages) == 1 - package = json.loads(packages[0].read_text(encoding="utf-8")) - assert package["metrics"][0]["name"] == "hermes.model_route.count" - assert package["metrics"][0]["value"] == 1 - assert (tmp_path / "atif" / "hermes-atif-s1.json").exists() - - -def test_real_binding_shares_plugin_configuration_across_two_profiles( - tmp_path, - monkeypatch, -): - relay = pytest.importorskip("nemo_relay") - if getattr(relay, "_native", None) is None: - pytest.skip("NeMo Relay native binding is unavailable on this platform") - plugin = _fresh_plugin(monkeypatch, relay) - original_initialize = relay.plugin.initialize - original_clear = relay.plugin.clear - original_clear() - initialize_calls = [] - clear_calls = 0 - - async def _initialize(config): - initialize_calls.append(config) - return await original_initialize(config) - - def _clear(): - nonlocal clear_calls - clear_calls += 1 - return original_clear() - - monkeypatch.setattr(relay.plugin, "initialize", _initialize) - monkeypatch.setattr(relay.plugin, "clear", _clear) - # This test exercises the bundled plugin's legacy configuration owner in - # isolation. Native and bundled-plugin ownership are intentionally not - # combined until their process-global lifetime models are unified. - monkeypatch.setattr( - relay_runtime._PLUGIN_CONFIGURATION, - "acquire", - lambda _owner, _relay: False, - ) - monkeypatch.setattr( - plugin, - "_load_settings", - lambda: plugin._Settings(plugins_config={"version": 1}), - ) - profile_a = str(tmp_path / "profile-a") - profile_b = str(tmp_path / "profile-b") - host_a = relay_runtime.RelayRuntime(relay=relay, profile_key=profile_a) - host_b = relay_runtime.RelayRuntime(relay=relay, profile_key=profile_b) - - try: - runtime_a = plugin._get_runtime(profile_key=profile_a, host=host_a) - runtime_b = plugin._get_runtime(profile_key=profile_b, host=host_b) - assert runtime_a is not None - assert runtime_b is not None - runtime_a.ensure_session({"session_id": "session-a"}) - runtime_b.ensure_session({"session_id": "session-b"}) - - assert initialize_calls == [{"version": 1}] - assert relay.plugin.report() is not None - - runtime_a.close_session({"session_id": "session-a"}) - - assert clear_calls == 0 - assert relay.plugin.report() is not None - assert runtime_b.host.get_session("session-b") is not None - - runtime_b.close_session({"session_id": "session-b"}) - - assert clear_calls == 1 - assert relay.plugin.report() is None - finally: - plugin.reset_for_tests() - host_a.shutdown() - host_b.shutdown() - original_clear() - - -def test_relay_tool_request_rewrite_precedes_hermes_authorization_boundary( - tmp_path, - monkeypatch, -): - from hermes_cli.middleware import apply_tool_request_middleware - - fake = _FakeNemoRelay() - plugin = _fresh_plugin(monkeypatch, fake) - _enable_dynamic_plugin(tmp_path, monkeypatch) - plugin.on_session_start(session_id="s1") - - result = apply_tool_request_middleware( - "fixture-tool", - {"value": 1}, - session_id="s1", - tool_call_id="tool-1", - ) - - assert result.payload == {"intercepted": True, "value": 1} - assert result.trace[0] == {"source": "nemo_relay"} diff --git a/website/docs/user-guide/features/built-in-plugins.md b/website/docs/user-guide/features/built-in-plugins.md index 77edbf7116..bb8f36e73c 100644 --- a/website/docs/user-guide/features/built-in-plugins.md +++ b/website/docs/user-guide/features/built-in-plugins.md @@ -58,7 +58,6 @@ The repo ships these bundled plugins under `plugins/`. All are opt-in — enable | `disk-cleanup` | hooks + slash command | Auto-track ephemeral files and clean them on session end | | `security-guidance` | hooks | Pattern-match dangerous code on `write_file`/`patch` and append a security warning (or block) — 25 rules (Apache-2.0 fork of Anthropic's `claude-plugins-official` patterns) | | `observability/langfuse` | hooks | Trace turns / LLM calls / tools to [Langfuse](https://langfuse.com) | -| `observability/nemo_relay` | hooks | Relay observability events (turns / LLM calls / tools) to an NVIDIA NeMo endpoint | | `teams_pipeline` | standalone | Microsoft Teams meeting pipeline — Graph-backed, transcript-first meeting summaries | | `spotify` | backend (7 tools) | Native Spotify playback, queue, search, playlists, albums, library | | `google_meet` | standalone | Join Meet calls, live-caption transcription, optional realtime duplex audio |