diff --git a/docs/middleware/README.md b/docs/middleware/README.md index 96304fa0e1..4a5c06f8cb 100644 --- a/docs/middleware/README.md +++ b/docs/middleware/README.md @@ -233,9 +233,8 @@ 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. -NeMo Relay execution middleware is installed through an explicitly selected -Relay `plugins.toml`; see -[Relay shared metrics](../observability/relay-shared-metrics.md). +For NeMo Relay adaptive execution middleware, see +[`plugins/observability/nemo_relay/README.md`](../../plugins/observability/nemo_relay/README.md). ## Safety Notes diff --git a/docs/observability/README.md b/docs/observability/README.md index f5ebea6177..915b0d5e2a 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 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). +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). diff --git a/docs/observability/monitoring.md b/docs/observability/monitoring.md index cba4ead1ff..f798bf2dbe 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 Hermes's native NeMo Relay SDK -integration and explicitly configured Relay subscribers or exporters. +capture is a separate plane served by the NeMo Relay integration +(`plugins/observability/nemo_relay/`) and its Hermes-owned subscribers. ## What gets exported diff --git a/docs/observability/relay-shared-metrics.md b/docs/observability/relay-shared-metrics.md index 3ae597b5d6..860d916e21 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 a Hermes observability plugin. 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 `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: 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.7.1 or later within the 0.7 release line. That +Hermes requires NeMo Relay 0.6.0 or later within the 0.6 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.7.1,<0.8` dependency range. The published package is built from +bounded `>=0.6.0,<0.7` 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,12 +41,9 @@ 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. -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. +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. 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 2c258ee8e8..6f7446620e 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 (``langfuse``), the directory - name, or the full path-derived key (``observability/langfuse``) and + Accepts either the bare manifest name (``nemo_relay``), the directory + name, or the full path-derived key (``observability/nemo_relay``) 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/langfuse``) included. + nested category plugins (e.g. ``observability/nemo_relay``) 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. "langfuse" -> - # "observability/langfuse"), but only when it resolves to exactly one + # 2. Fall back to a bare leaf-name match (e.g. "nemo_relay" -> + # "observability/nemo_relay"), 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 new file mode 100644 index 0000000000..85f7d03903 --- /dev/null +++ b/plugins/observability/nemo_relay/README.md @@ -0,0 +1,602 @@ +# 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 new file mode 100644 index 0000000000..13218ed5e4 --- /dev/null +++ b/plugins/observability/nemo_relay/__init__.py @@ -0,0 +1,1023 @@ +"""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 new file mode 100644 index 0000000000..046d5d0d85 --- /dev/null +++ b/plugins/observability/nemo_relay/plugin.yaml @@ -0,0 +1,15 @@ +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 5a781d0b32..0f8525eaac 100644 --- a/scripts/toolperf_abeval/README.md +++ b/scripts/toolperf_abeval/README.md @@ -40,11 +40,9 @@ 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 01fa9a355b..176685ecd3 100644 --- a/scripts/toolperf_abeval/ab_eval.py +++ b/scripts/toolperf_abeval/ab_eval.py @@ -165,34 +165,14 @@ 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_PLUGINS_TOML": str(relay_config), + "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", }) 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 38e1b57a7c..a964626e23 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/trace_sink`) to the canonical +OR a full path-derived key (e.g. `observability/nemo_relay`) 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", "trace_sink", { - "name": "trace_sink", "version": "1.0.0", "description": "trace sink" + _make_category_plugin(tmp_path, "observability", "nemo_relay", { + "name": "nemo_relay", "version": "1.0.0", "description": "relay obs" }) _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/trace_sink") == "observability/trace_sink" + assert _resolve_plugin_key("observability/nemo_relay") == "observability/nemo_relay" @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("trace_sink", allow_tool_override=False) # bare name + cmd_enable("nemo_relay", 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/trace_sink" in saved - assert "trace_sink" not in saved or "observability/trace_sink" in saved + assert "observability/nemo_relay" in saved + assert "nemo_relay" not in saved or "observability/nemo_relay" 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 new file mode 100644 index 0000000000..343bd76965 --- /dev/null +++ b/tests/plugins/test_nemo_relay_plugin.py @@ -0,0 +1,490 @@ +"""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 bb8f36e73c..77edbf7116 100644 --- a/website/docs/user-guide/features/built-in-plugins.md +++ b/website/docs/user-guide/features/built-in-plugins.md @@ -58,6 +58,7 @@ 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 |