From 64faff6768ffac903b84ab6abd94e1ba3b3cc58b Mon Sep 17 00:00:00 2001 From: Alex Fournier Date: Sun, 19 Jul 2026 08:54:49 -0400 Subject: [PATCH] refactor(observability): move Relay runtime ownership into core Signed-off-by: Alex Fournier --- MANIFEST.in | 3 + docs/observability/README.md | 4 + docs/observability/relay-shared-metrics.md | 72 ++ hermes_cli/observability/__init__.py | 51 ++ hermes_cli/observability/relay_runtime.py | 372 +++++++++ .../observability/relay_shared_metrics.py | 342 ++++++++ .../hermes.shared_metrics.v1.schema.json | 0 .../observability}/shared_metrics.py | 0 .../observability}/shared_metrics_contract.py | 29 +- .../shared_metrics_subscriber.py | 0 hermes_cli/plugins.py | 33 +- plugins/observability/nemo_relay/README.md | 84 +- plugins/observability/nemo_relay/__init__.py | 770 +++++------------- pyproject.toml | 8 +- scripts/smoke_nemo_relay_shared_metrics.py | 34 +- .../test_relay_shared_metrics.py} | 14 +- .../test_relay_shared_metrics_runtime.py | 463 +++++++++++ tests/plugins/test_nemo_relay_plugin.py | 721 +++++----------- tests/test_project_metadata.py | 52 +- uv.lock | 8 +- 20 files changed, 1826 insertions(+), 1234 deletions(-) create mode 100644 docs/observability/relay-shared-metrics.md create mode 100644 hermes_cli/observability/__init__.py create mode 100644 hermes_cli/observability/relay_runtime.py create mode 100644 hermes_cli/observability/relay_shared_metrics.py rename {plugins/observability/nemo_relay => hermes_cli/observability}/schemas/hermes.shared_metrics.v1.schema.json (100%) rename {plugins/observability/nemo_relay => hermes_cli/observability}/shared_metrics.py (100%) rename {plugins/observability/nemo_relay => hermes_cli/observability}/shared_metrics_contract.py (91%) rename {plugins/observability/nemo_relay => hermes_cli/observability}/shared_metrics_subscriber.py (100%) rename tests/{plugins/test_nemo_relay_shared_metrics.py => hermes_cli/test_relay_shared_metrics.py} (97%) create mode 100644 tests/hermes_cli/test_relay_shared_metrics_runtime.py diff --git a/MANIFEST.in b/MANIFEST.in index 159c215ff6b0..c5c33210632e 100644 --- a/MANIFEST.in +++ b/MANIFEST.in @@ -7,6 +7,9 @@ graft locales # built from the sdist (e.g. Homebrew, downstream packagers). package-data # below covers the wheel; this covers the sdist. See #34034 / #28149. recursive-include plugins plugin.yaml plugin.yml +# The built-in shared-metrics package validates exported JSON against this +# closed schema. Include it when downstream packagers build from the sdist. +recursive-include hermes_cli/observability/schemas *.json # Gateway assets include images plus YAML catalogs such as status_phrases.yaml. recursive-include gateway/assets * global-exclude __pycache__ diff --git a/docs/observability/README.md b/docs/observability/README.md index 9040929ca489..7ecbd2e7bd6b 100644 --- a/docs/observability/README.md +++ b/docs/observability/README.md @@ -14,6 +14,10 @@ Behavior-changing request or execution wrappers are outside this observer contract. Observer hooks should report what happened; they should not replace provider requests, tool arguments, or execution callbacks. +Hermes also has a first-party NeMo Relay shared-metrics path. It uses these +lifecycle boundaries directly and does not require enabling an observability +plugin. See [Relay shared metrics](relay-shared-metrics.md). + ## Contract Plugins register observer callbacks from `register(ctx)`: diff --git a/docs/observability/relay-shared-metrics.md b/docs/observability/relay-shared-metrics.md new file mode 100644 index 000000000000..6c85c84cda43 --- /dev/null +++ b/docs/observability/relay-shared-metrics.md @@ -0,0 +1,72 @@ +# NeMo Relay Shared Metrics + +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, but Relay-backed instrumentation is unavailable there. + +Collection remains off unless Hermes policy enables it: + +```yaml +telemetry: + shared_metrics: + enabled: true +``` + +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 +`hermes_cli.observability.relay_runtime` to obtain the shared session handle or +run Relay scope, LLM, tool, and mark APIs in that session context. New product +marks do not require Hermes plugin registration. Shared-metrics marks must +still contain only fields approved by the versioned allowlist; the hard +dependency does not change the collection or privacy policy. + +## Current Slice + +The first vertical slice records one logical model-call counter: + +```text +Hermes API hooks + -> Relay session and LLM lifecycle + -> Hermes shared-metrics subscriber + -> SQLite counter + -> immutable JSON delta package +``` + +Hermes sends an empty `LLMRequest` into this metrics lifecycle. The terminal +event contains only bounded model family, provider family, locality, call role, +and outcome values. Prompts, responses, exact model IDs, endpoints, errors, +session IDs, task IDs, and request IDs are not included in the metrics event or +package. + +Local state is written under: + +```text +$HERMES_HOME/telemetry/shared_metrics/metrics.sqlite3 +$HERMES_HOME/telemetry/shared_metrics/outbox/*.json +``` + +The database keeps transactional aggregate and package-outbox state. Package +files are immutable delta documents that conform to a closed JSON schema and +are written with atomic replacement. + +## Smoke Test + +Run a real Hermes CLI turn against the deterministic local model server: + +```bash +./.venv/bin/python scripts/smoke_nemo_relay_shared_metrics.py +``` + +The script uses the installed `nemo-relay` dependency by default. Pass +`--relay-python ../nemo-relay/python` only when testing a locally built Relay +binding. + +The smoke verifies the model request reached the local server, one counter was +stored, one package was exported, and prompt, response, and exact-model +canaries are absent from the package. diff --git a/hermes_cli/observability/__init__.py b/hermes_cli/observability/__init__.py new file mode 100644 index 000000000000..ce01554d93e5 --- /dev/null +++ b/hermes_cli/observability/__init__.py @@ -0,0 +1,51 @@ +"""First-party Hermes observability integrations.""" + +from __future__ import annotations + +import logging +from typing import Any + +logger = logging.getLogger(__name__) + + +def prepare_lifecycle(hook_name: str, **kwargs: Any) -> None: + """Prepare subscribers that must observe a session's start event.""" + from . import relay_runtime, relay_shared_metrics + + if hook_name in relay_runtime.SESSION_START_HOOKS: + try: + relay_shared_metrics.prepare_session_start() + except Exception: + logger.warning("Built-in observability preparation failed", exc_info=True) + + +def observe_lifecycle(hook_name: str, **kwargs: Any) -> None: + """Dispatch a Hermes lifecycle event to built-in observability features.""" + from . import relay_runtime, relay_shared_metrics + + # Session-start plugin callbacks register optional per-session subscribers + # before this completion step opens the shared core scope. On teardown, + # metrics finish child LLM scopes before the neutral host closes the owner. + if hook_name not in relay_runtime.SESSION_CLOSE_HOOKS: + _safe_observe(relay_runtime.observe_lifecycle, hook_name, kwargs) + _safe_observe(relay_shared_metrics.observe_lifecycle, hook_name, kwargs) + if hook_name in relay_runtime.SESSION_CLOSE_HOOKS: + _safe_observe(relay_runtime.observe_lifecycle, hook_name, kwargs) + + +def handles_hook(hook_name: str) -> bool: + """Return whether any built-in observability feature handles a hook.""" + from . import relay_runtime, relay_shared_metrics + + return relay_runtime.handles_hook(hook_name) or relay_shared_metrics.handles_hook( + hook_name + ) + + +def _safe_observe(callback: Any, hook_name: str, kwargs: dict[str, Any]) -> None: + try: + callback(hook_name, **kwargs) + except Exception: + logger.warning( + "Built-in observability hook failed: %s", hook_name, exc_info=True + ) diff --git a/hermes_cli/observability/relay_runtime.py b/hermes_cli/observability/relay_runtime.py new file mode 100644 index 000000000000..d52b44d07cb1 --- /dev/null +++ b/hermes_cli/observability/relay_runtime.py @@ -0,0 +1,372 @@ +"""Process-wide NeMo Relay runtime owned by Hermes core.""" + +from __future__ import annotations + +import atexit +import contextvars +import importlib +import logging +import threading +from dataclasses import dataclass, field +from typing import Any, Callable + +logger = logging.getLogger(__name__) + +SESSION_SCOPE = "hermes.session" +RUNTIME_SCHEMA_KEY = "hermes.relay.schema_version" +RUNTIME_SCHEMA_VERSION = "hermes.relay.runtime.v1" + +SESSION_START_HOOKS = frozenset({"on_session_start"}) +SESSION_CLOSE_HOOKS = frozenset({"on_session_finalize", "on_session_reset"}) +SUBAGENT_START_HOOKS = frozenset({"subagent_start"}) +SUBAGENT_STOP_HOOKS = frozenset({"subagent_stop"}) +HANDLED_HOOKS = ( + SESSION_START_HOOKS + | SESSION_CLOSE_HOOKS + | SUBAGENT_START_HOOKS + | SUBAGENT_STOP_HOOKS +) + +_RUNTIME_FAILED = object() +_RUNTIME: RelayRuntime | object | None = None +_RUNTIME_LOCK = threading.RLock() + + +@dataclass +class RelaySession: + """One isolated Relay scope stack owned by a Hermes session.""" + + session_id: str + parent_session_id: str = "" + lock: threading.RLock = field(default_factory=threading.RLock, repr=False) + closing: bool = False + handle: Any = None + context: contextvars.Context | None = None + + +class RelayRuntime: + """Own Relay session scopes independently of any exporter or plugin.""" + + def __init__(self, relay: Any = None) -> None: + self.relay = relay or _load_nemo_relay() + self._sessions_lock = threading.RLock() + self._sessions: dict[str, RelaySession] = {} + self._subagent_parents: dict[str, str] = {} + self._shutdown_registered = True + atexit.register(self.shutdown) + + def ensure_session( + self, + event: dict[str, Any], + *, + data: Any = None, + metadata: dict[str, Any] | None = None, + ) -> RelaySession | None: + """Return the existing session scope or create it once.""" + session_id = _session_id(event) + if not session_id: + return None + with self._sessions_lock: + session = self._sessions.get(session_id) + if session is None: + parent_session_id = self._subagent_parents.get(session_id, "") + session = RelaySession( + session_id=session_id, + parent_session_id=parent_session_id, + ) + self._sessions[session_id] = session + with session.lock: + if session.closing: + return None + if session.handle is None: + parent_handle = None + scope_metadata = { + **(metadata or {}), + RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION, + } + if session.parent_session_id: + parent = self.ensure_session({ + "session_id": session.parent_session_id + }) + if parent is not None: + parent_handle = parent.handle + scope_metadata["nemo_relay_scope_role"] = "subagent" + context = contextvars.Context() + try: + session.handle = context.run( + self.relay.scope.push, + SESSION_SCOPE, + self.relay.ScopeType.Agent, + handle=parent_handle, + data=data, + input={}, + metadata=scope_metadata, + ) + except Exception: + session.context = None + raise + session.context = context + return session + + def register_subagent(self, event: dict[str, Any]) -> None: + """Record the parent used when a delegated Hermes session starts.""" + parent_session_id = str(event.get("parent_session_id") or "") + child_session_id = str(event.get("child_session_id") or "") + if ( + not parent_session_id + or not child_session_id + or parent_session_id == child_session_id + ): + return + self.ensure_session({"session_id": parent_session_id}) + with self._sessions_lock: + self._subagent_parents[child_session_id] = parent_session_id + + def unregister_subagent(self, event: dict[str, Any]) -> None: + """Forget a delegated-session relationship after its terminal hook.""" + child_session_id = str(event.get("child_session_id") or "") + if not child_session_id: + return + with self._sessions_lock: + self._subagent_parents.pop(child_session_id, None) + + def get_session(self, session_id: str) -> RelaySession | None: + """Return an active Hermes Relay session without creating one.""" + with self._sessions_lock: + session = self._sessions.get(str(session_id or "")) + if session is None: + return None + with session.lock: + return None if session.closing else session + + def get_session_handle(self, session_id: str) -> Any: + """Return the Relay parent handle for a Hermes session, if active.""" + session = self.get_session(session_id) + return None if session is None else session.handle + + def run_in_session( + self, + session: RelaySession, + callback: Callable[..., Any], + *args: Any, + allow_closing: bool = False, + **kwargs: Any, + ) -> Any: + """Run a Relay operation against a session's isolated scope stack.""" + with session.lock: + if session.closing and not allow_closing: + raise RuntimeError("Hermes Relay session is closing") + if session.context is None or session.handle is None: + raise RuntimeError("Hermes Relay session context is unavailable") + + def invoke() -> Any: + self.relay.get_scope_stack() + return callback(*args, **kwargs) + + # A copy permits a helper called by an existing Relay callback to + # re-enter the same logical session without re-entering Context. + return session.context.copy().run(invoke) + + def emit_mark( + self, + name: str, + event: dict[str, Any], + *, + data: Any = None, + metadata: Any = None, + ) -> bool: + """Emit a mark parented to the Hermes session identified by ``event``.""" + session = self.ensure_session(event) + if session is None: + return False + self.run_in_session( + session, + self.relay.scope.event, + name, + handle=session.handle, + data=data, + metadata=metadata, + ) + return True + + def close_session(self, event: dict[str, Any]) -> None: + """Close one session scope and remove it from the core registry.""" + session_id = _session_id(event) + with self._sessions_lock: + session = self._sessions.get(session_id) + if session is None: + return + failures: list[str] = [] + with session.lock: + if session.closing: + return + session.closing = True + if session.handle is not None: + try: + self.run_in_session( + session, + self.relay.scope.pop, + session.handle, + output={}, + metadata={RUNTIME_SCHEMA_KEY: RUNTIME_SCHEMA_VERSION}, + allow_closing=True, + ) + except Exception as exc: + failures.append(f"session scope close failed: {exc}") + try: + self.relay.subscribers.flush() + except Exception as exc: + failures.append(f"subscriber flush failed: {exc}") + with self._sessions_lock: + if self._sessions.get(session_id) is session: + self._sessions.pop(session_id, None) + self._subagent_parents.pop(session_id, None) + if failures: + logger.warning( + "Hermes Relay session %s closed with errors: %s", + session_id, + "; ".join(failures), + ) + + def shutdown(self) -> None: + """Close all core-owned Relay session scopes.""" + with self._sessions_lock: + session_ids = list(self._sessions) + for session_id in session_ids: + self._safe(self.close_session, {"session_id": session_id}) + if self._shutdown_registered: + try: + atexit.unregister(self.shutdown) + except Exception: + pass + self._shutdown_registered = False + + @staticmethod + def _safe(callback: Callable[..., Any], *args: Any, **kwargs: Any) -> Any: + try: + return callback(*args, **kwargs) + except Exception: + logger.warning("Hermes Relay runtime operation failed", exc_info=True) + return None + + +def handles_hook(hook_name: str) -> bool: + """Return whether the core Relay host consumes this lifecycle hook.""" + return hook_name in HANDLED_HOOKS + + +def observe_lifecycle(hook_name: str, **kwargs: Any) -> None: + """Apply session lifecycle events to the core Relay host.""" + if not handles_hook(hook_name): + return + # Session hooks do not activate Relay by themselves. A direct core + # producer or an enabled built-in consumer creates the host lazily, after + # which these hooks keep its session lifetime correct. + runtime = get_runtime(create=False) + if runtime is None: + return + try: + if hook_name in SESSION_START_HOOKS: + runtime.ensure_session(kwargs) + elif hook_name in SESSION_CLOSE_HOOKS: + runtime.close_session(kwargs) + elif hook_name in SUBAGENT_START_HOOKS: + runtime.register_subagent(kwargs) + else: + runtime.unregister_subagent(kwargs) + except Exception: + logger.warning("Hermes Relay lifecycle failed: %s", hook_name, exc_info=True) + + +def emit_mark( + name: str, + *, + session_id: str, + data: Any = None, + metadata: Any = None, +) -> bool: + """Emit a fail-open Relay mark under a Hermes session.""" + runtime = get_runtime() + if runtime is None: + return False + try: + return runtime.emit_mark( + name, + {"session_id": session_id}, + data=data, + metadata=metadata, + ) + except Exception: + logger.warning("Hermes Relay mark failed: %s", name, exc_info=True) + return False + + +def ensure_session(*, session_id: str, **context: Any) -> RelaySession | None: + """Create or return the shared Relay session used by Hermes core.""" + runtime = get_runtime() + if runtime is None: + return None + try: + return runtime.ensure_session({"session_id": session_id, **context}) + except Exception: + logger.warning("Hermes Relay session initialization failed", exc_info=True) + return None + + +def run_in_session( + session_id: str, + callback: Callable[..., Any], + *args: Any, + **kwargs: Any, +) -> Any: + """Run a scope, LLM, or tool API against a shared Hermes session.""" + runtime = get_runtime() + if runtime is None: + raise RuntimeError("Hermes Relay runtime is unavailable") + session = runtime.get_session(session_id) + if session is None: + session = runtime.ensure_session({"session_id": session_id}) + if session is None: + raise RuntimeError("Hermes Relay session is unavailable") + return runtime.run_in_session(session, callback, *args, **kwargs) + + +def get_session_handle(session_id: str) -> Any: + """Return the shared Relay handle for direct core instrumentation.""" + runtime = get_runtime(create=False) + return None if runtime is None else runtime.get_session_handle(session_id) + + +def get_runtime(*, create: bool = True) -> RelayRuntime | None: + """Return the process-wide Hermes Relay host.""" + global _RUNTIME + with _RUNTIME_LOCK: + if isinstance(_RUNTIME, RelayRuntime): + return _RUNTIME + if _RUNTIME is _RUNTIME_FAILED or not create: + return None + try: + _RUNTIME = RelayRuntime() + except Exception: + logger.warning("Hermes Relay runtime initialization failed", exc_info=True) + _RUNTIME = _RUNTIME_FAILED + return None + return _RUNTIME + + +def _load_nemo_relay() -> Any: + """Load the binding only when a producer or consumer needs Relay.""" + return importlib.import_module("nemo_relay") + + +def _session_id(event: dict[str, Any]) -> str: + return str(event.get("session_id") or "") + + +def _reset_for_tests() -> None: + """Reset process-global core Relay state for isolated tests.""" + global _RUNTIME + with _RUNTIME_LOCK: + if isinstance(_RUNTIME, RelayRuntime): + _RUNTIME.shutdown() + _RUNTIME = None diff --git a/hermes_cli/observability/relay_shared_metrics.py b/hermes_cli/observability/relay_shared_metrics.py new file mode 100644 index 000000000000..c8306966f1e7 --- /dev/null +++ b/hermes_cli/observability/relay_shared_metrics.py @@ -0,0 +1,342 @@ +"""Direct NeMo Relay integration for Hermes shared client metrics.""" + +from __future__ import annotations + +import atexit +import logging +import threading +from dataclasses import dataclass, field +from functools import lru_cache +from typing import Any, Callable + +from hermes_cli import __version__ + +from . import relay_runtime +from .shared_metrics import SharedMetricsStore +from .shared_metrics_contract import ( + MODEL_CALL_SCOPE, + SCHEMA_KEY, + SCHEMA_VERSION, + SUBSCRIBER_NAME, + model_call_fields, + model_call_outcome, +) +from .shared_metrics_subscriber import SharedMetricsSubscriber + +logger = logging.getLogger(__name__) + +HANDLED_HOOKS = frozenset({ + "on_session_start", + "on_session_end", + "on_session_finalize", + "on_session_reset", + "pre_api_request", + "post_api_request", + "api_request_error", +}) + +_RUNTIME_FAILED = object() +_RUNTIME: _Runtime | object | None = None +_RUNTIME_LOCK = threading.RLock() + + +@dataclass +class _ModelCall: + handle: Any + task_id: str + fields: dict[str, str] + + +@dataclass +class _MetricsSession: + session_id: str + relay_session: relay_runtime.RelaySession + lock: threading.RLock = field(default_factory=threading.RLock, repr=False) + closing: bool = False + model_calls: dict[str, _ModelCall] = field(default_factory=dict) + + +class _Runtime: + """Own shared-metrics state layered on the Hermes core Relay host.""" + + def __init__(self, host: relay_runtime.RelayRuntime | None = None) -> None: + resolved_host = host or relay_runtime.get_runtime() + if resolved_host is None: + raise RuntimeError("Hermes core Relay runtime is unavailable") + self.host: relay_runtime.RelayRuntime = resolved_host + self.relay = self.host.relay + self._sessions_lock = threading.RLock() + self._sessions: dict[str, _MetricsSession] = {} + self.subscriber = SharedMetricsSubscriber(SharedMetricsStore(), __version__) + self.relay.subscribers.register(SUBSCRIBER_NAME, self.subscriber) + self._registered = True + atexit.register(self.shutdown) + + def ensure_session(self, event: dict[str, Any]) -> _MetricsSession | None: + session_id = str(event.get("session_id") or "") + if not session_id: + return None + relay_session = self.host.ensure_session(event) + if relay_session is None: + return None + with self._sessions_lock: + session = self._sessions.get(session_id) + if session is None: + session = _MetricsSession( + session_id=session_id, + relay_session=relay_session, + ) + self._sessions[session_id] = session + with session.lock: + if session.closing: + return None + return session + + def _run_in_session( + self, + session: _MetricsSession, + callback: Callable[..., Any], + *args: Any, + **kwargs: Any, + ) -> Any: + return self.host.run_in_session( + session.relay_session, + callback, + *args, + **kwargs, + ) + + def start_model_call(self, event: dict[str, Any]) -> None: + session = self.ensure_session(event) + if session is None: + return + request_id = str(event.get("api_request_id") or "") + if not request_id: + return + fields = model_call_fields(event) + model_family = fields["model_family"] + with session.lock: + if session.closing: + return + existing = session.model_calls.get(request_id) + if existing is not None: + existing.fields = fields + return + handle = self._run_in_session( + session, + self.relay.llm.call, + MODEL_CALL_SCOPE, + self.relay.LLMRequest({}, {}), + handle=session.relay_session.handle, + metadata={SCHEMA_KEY: SCHEMA_VERSION}, + model_name=model_family, + ) + session.model_calls[request_id] = _ModelCall( + handle=handle, + task_id=str(event.get("task_id") or ""), + fields=fields, + ) + + def end_model_call(self, event: dict[str, Any], outcome: str | None = None) -> None: + session = self._session(event) + if session is None: + return + request_id = str(event.get("api_request_id") or "") + with session.lock: + if session.closing: + return + model_call = session.model_calls.get(request_id) + if model_call is None: + return + fields = model_call_fields(event) + model_call.fields = fields + self._finish_model_call( + session, + request_id, + outcome or model_call_outcome(event), + ) + + def end_pending_model_calls(self, event: dict[str, Any]) -> None: + session = self._session(event) + if session is None: + return + with session.lock: + if session.closing: + return + self._end_pending_model_calls(session, event) + + def close_session(self, event: dict[str, Any]) -> None: + session = self._session(event) + if session is None: + return + failures: list[str] = [] + with session.lock: + if session.closing: + return + session.closing = True + self._end_pending_model_calls(session, event) + try: + self.relay.subscribers.flush() + except Exception as exc: + failures.append(f"subscriber flush failed: {exc}") + self._export() + with self._sessions_lock: + if self._sessions.get(session.session_id) is session: + self._sessions.pop(session.session_id, None) + if failures: + logger.warning( + "Hermes shared-metrics session %s closed with errors: %s", + session.session_id, + "; ".join(failures), + ) + + def shutdown(self) -> None: + with self._sessions_lock: + session_ids = list(self._sessions) + for session_id in session_ids: + self._safe(self.close_session, {"session_id": session_id}) + if not self._registered: + return + self._safe(self.relay.subscribers.flush) + self._export() + self._safe(self.relay.subscribers.deregister, SUBSCRIBER_NAME) + self._registered = False + try: + atexit.unregister(self.shutdown) + except Exception: + pass + + def _session(self, event: dict[str, Any]) -> _MetricsSession | None: + session_id = str(event.get("session_id") or "") + with self._sessions_lock: + return self._sessions.get(session_id) + + def _finish_model_call( + self, + session: _MetricsSession, + request_id: str, + outcome: str, + ) -> None: + model_call = session.model_calls.pop(request_id, None) + if model_call is None: + return + try: + self._run_in_session( + session, + self.relay.llm.call_end, + model_call.handle, + {**model_call.fields, "outcome": outcome}, + metadata={SCHEMA_KEY: SCHEMA_VERSION}, + ) + except Exception: + logger.warning( + "Hermes shared-metrics model call close failed", exc_info=True + ) + + def _end_pending_model_calls( + self, + session: _MetricsSession, + event: dict[str, Any], + ) -> None: + task_id = str(event.get("task_id") or "") + request_ids = [ + request_id + for request_id, model_call in session.model_calls.items() + if not task_id or model_call.task_id == task_id + ] + outcome = "cancelled" if event.get("interrupted") else "failed" + for request_id in request_ids: + self._finish_model_call(session, request_id, outcome) + + def _export(self) -> None: + self._safe(self.subscriber.store.create_and_export_package) + + @staticmethod + def _safe(callback: Callable[..., Any], *args: Any, **kwargs: Any) -> Any: + try: + return callback(*args, **kwargs) + except Exception: + logger.warning("Hermes shared metrics operation failed", exc_info=True) + return None + + +@lru_cache(maxsize=1) +def enabled() -> bool: + """Return the process-lifetime Hermes shared-metrics policy.""" + try: + from hermes_cli.config import load_config_readonly + + config = load_config_readonly() or {} + except Exception: + logger.debug("Unable to read Hermes shared-metrics policy", exc_info=True) + return False + if not isinstance(config, dict): + return False + telemetry = config.get("telemetry") + if not isinstance(telemetry, dict): + return False + shared_metrics = telemetry.get("shared_metrics") + return isinstance(shared_metrics, dict) and shared_metrics.get("enabled") is True + + +def handles_hook(hook_name: str) -> bool: + return hook_name in HANDLED_HOOKS and enabled() + + +def observe_lifecycle(hook_name: str, **kwargs: Any) -> None: + """Project one Hermes lifecycle event into the core Relay integration.""" + if not handles_hook(hook_name): + return + runtime = _get_runtime() + if runtime is None: + return + try: + if hook_name == "on_session_start": + runtime.ensure_session(kwargs) + elif hook_name == "pre_api_request": + runtime.start_model_call(kwargs) + elif hook_name == "post_api_request": + runtime.end_model_call(kwargs, "success") + elif hook_name == "api_request_error": + if kwargs.get("retryable") is False: + runtime.end_model_call(kwargs, "failed") + elif hook_name == "on_session_end": + runtime.end_pending_model_calls(kwargs) + elif hook_name in {"on_session_finalize", "on_session_reset"}: + runtime.close_session(kwargs) + except Exception: + logger.warning( + "Hermes shared metrics hook failed: %s", hook_name, exc_info=True + ) + + +def prepare_session_start() -> None: + """Register the subscriber before any producer opens the session scope.""" + if enabled(): + _get_runtime() + + +def _get_runtime() -> _Runtime | None: + global _RUNTIME + with _RUNTIME_LOCK: + if isinstance(_RUNTIME, _Runtime): + return _RUNTIME + if _RUNTIME is _RUNTIME_FAILED: + return None + try: + _RUNTIME = _Runtime() + except Exception: + logger.warning("Hermes shared metrics initialization failed", exc_info=True) + _RUNTIME = _RUNTIME_FAILED + return None + return _RUNTIME + + +def _reset_for_tests() -> None: + """Reset process-global state for isolated tests.""" + global _RUNTIME + with _RUNTIME_LOCK: + if isinstance(_RUNTIME, _Runtime): + _RUNTIME.shutdown() + _RUNTIME = None + enabled.cache_clear() diff --git a/plugins/observability/nemo_relay/schemas/hermes.shared_metrics.v1.schema.json b/hermes_cli/observability/schemas/hermes.shared_metrics.v1.schema.json similarity index 100% rename from plugins/observability/nemo_relay/schemas/hermes.shared_metrics.v1.schema.json rename to hermes_cli/observability/schemas/hermes.shared_metrics.v1.schema.json diff --git a/plugins/observability/nemo_relay/shared_metrics.py b/hermes_cli/observability/shared_metrics.py similarity index 100% rename from plugins/observability/nemo_relay/shared_metrics.py rename to hermes_cli/observability/shared_metrics.py diff --git a/plugins/observability/nemo_relay/shared_metrics_contract.py b/hermes_cli/observability/shared_metrics_contract.py similarity index 91% rename from plugins/observability/nemo_relay/shared_metrics_contract.py rename to hermes_cli/observability/shared_metrics_contract.py index d7937a6483ab..7e3c5d336e83 100644 --- a/plugins/observability/nemo_relay/shared_metrics_contract.py +++ b/hermes_cli/observability/shared_metrics_contract.py @@ -8,12 +8,11 @@ from typing import Any SCHEMA_KEY = "hermes.metrics.schema_version" SCHEMA_VERSION = "hermes.metrics.event.v1" -SESSION_SCOPE = "hermes.session" MODEL_CALL_SCOPE = "hermes.model_call" SUBSCRIBER_NAME = "hermes.nemo_relay.shared_metrics" PRIMARY_MODEL_CALL_ROLE = "primary" -EXECUTION_SURFACES = frozenset({ +EXECUTION_SURFACES: frozenset[str] = frozenset({ "api", "batch", "cli", @@ -25,14 +24,20 @@ EXECUTION_SURFACES = frozenset({ "other", "unknown", }) -PROVIDER_FAMILIES = frozenset({"aggregator", "custom", "direct", "local", "unknown"}) -MODEL_LOCALITIES = frozenset({"local", "remote", "unknown"}) -MODEL_OUTCOMES = frozenset({"cancelled", "failed", "success"}) +PROVIDER_FAMILIES: frozenset[str] = frozenset({ + "aggregator", + "custom", + "direct", + "local", + "unknown", +}) +MODEL_LOCALITIES: frozenset[str] = frozenset({"local", "remote", "unknown"}) +MODEL_OUTCOMES: frozenset[str] = frozenset({"cancelled", "failed", "success"}) # Shared metrics use an explicit family allowlist rather than raw model IDs or # dynamically sourced catalog values. The latter would make the exported schema # drift independently of this contract. -MODEL_FAMILIES = frozenset({ +MODEL_FAMILIES: frozenset[str] = frozenset({ "claude", "deepseek", "gemini", @@ -60,7 +65,11 @@ _MODEL_FAMILY_PATTERN = re.compile( r"(?:^|[/_.:-])(" + "|".join( re.escape(family) - for family in sorted(MODEL_FAMILIES - {"unknown"}, key=len, reverse=True) + for family in sorted( + MODEL_FAMILIES - {"unknown"}, + key=lambda value: len(value), + reverse=True, + ) ) + r")(?=$|[/_.:-]|\d)" ) @@ -109,6 +118,7 @@ def model_call_dimensions(event: Any) -> dict[str, str] | None: expected_fields = { "call_role", "locality", + "model_family", "outcome", "provider_family", } @@ -117,6 +127,7 @@ def model_call_dimensions(event: Any) -> dict[str, str] | None: if ( data.get("call_role") != PRIMARY_MODEL_CALL_ROLE or data.get("locality") not in MODEL_LOCALITIES + or data.get("model_family") not in MODEL_FAMILIES or data.get("outcome") not in MODEL_OUTCOMES or data.get("provider_family") not in PROVIDER_FAMILIES ): @@ -124,7 +135,7 @@ def model_call_dimensions(event: Any) -> dict[str, str] | None: return { "call_role": PRIMARY_MODEL_CALL_ROLE, "locality": data["locality"], - "model_family": event_model_family, + "model_family": data["model_family"], "outcome": data["outcome"], "provider_family": data["provider_family"], } @@ -240,7 +251,7 @@ def model_family(kwargs: dict[str, Any]) -> str: declared_family = str(kwargs.get("model_family") or "").strip().lower() if declared_family in MODEL_FAMILIES - {"unknown"}: return declared_family - model = str(kwargs.get("model") or "").lower() + model = str(kwargs.get("response_model") or kwargs.get("model") or "").lower() match = _MODEL_FAMILY_PATTERN.search(model) return match.group(1) if match is not None else "unknown" diff --git a/plugins/observability/nemo_relay/shared_metrics_subscriber.py b/hermes_cli/observability/shared_metrics_subscriber.py similarity index 100% rename from plugins/observability/nemo_relay/shared_metrics_subscriber.py rename to hermes_cli/observability/shared_metrics_subscriber.py diff --git a/hermes_cli/plugins.py b/hermes_cli/plugins.py index 6ca393fca53c..9b1d24880812 100644 --- a/hermes_cli/plugins.py +++ b/hermes_cli/plugins.py @@ -2047,10 +2047,32 @@ def discover_plugins(force: bool = False) -> None: def invoke_hook(hook_name: str, **kwargs: Any) -> List[Any]: - """Invoke a lifecycle hook on all loaded plugins. + """Invoke a lifecycle hook on built-in observers and loaded plugins. Returns a list of non-``None`` return values from plugin callbacks. """ + if hook_name == "on_session_start": + try: + from hermes_cli.observability import prepare_lifecycle + + prepare_lifecycle(hook_name, **kwargs) + except Exception: + logger.warning("Built-in observability preparation failed", exc_info=True) + results = get_plugin_manager().invoke_hook(hook_name, **kwargs) + try: + from hermes_cli.observability import observe_lifecycle + + observe_lifecycle(hook_name, **kwargs) + except Exception: + logger.warning("Built-in observability hook failed", exc_info=True) + return results + + try: + from hermes_cli.observability import observe_lifecycle + + observe_lifecycle(hook_name, **kwargs) + except Exception: + logger.warning("Built-in observability hook failed", exc_info=True) return get_plugin_manager().invoke_hook(hook_name, **kwargs) @@ -2072,7 +2094,14 @@ def has_middleware(kind: str) -> bool: def has_hook(hook_name: str) -> bool: - """Return True when a hook has registered callbacks.""" + """Return True when a built-in observer or plugin handles a hook.""" + try: + from hermes_cli.observability import handles_hook + + if handles_hook(hook_name): + return True + except Exception: + logger.warning("Unable to inspect built-in shared-metrics hooks", exc_info=True) return get_plugin_manager().has_hook(hook_name) diff --git a/plugins/observability/nemo_relay/README.md b/plugins/observability/nemo_relay/README.md index 60db027072e8..52d31bb78903 100644 --- a/plugins/observability/nemo_relay/README.md +++ b/plugins/observability/nemo_relay/README.md @@ -73,89 +73,24 @@ checkout that contains this plugin. A globally installed older CLI will not see new bundled plugins from your working tree. ```bash -uv sync --extra nemo-relay +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, then install the official NeMo Relay runtime extra: +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 -python -m pip install "nemo-relay>=0.5.0,<0.6.0" hermes plugins enable observability/nemo_relay ``` -The plugin fails open when `nemo-relay` is not installed. Install a supported -NeMo Relay 0.5.x distribution: - -```bash -pip install "nemo-relay>=0.5.0,<0.6.0" -``` - -## Shared Metrics Proof Mode - -The Phase 1 telemetry proof adds a separately gated metrics-only mode. Enable -the bundled plugin and the mode in the same `config.yaml`: - -```yaml -plugins: - enabled: - - observability/nemo_relay - entries: - observability/nemo_relay: - shared_metrics: - enabled: true -``` - -The proof reads this setting when the plugin runtime initializes. Restart a -long-running Hermes agent or gateway after changing it. Live revocation and -local-state reset controls are follow-on requirements before production -rollout. - -This first vertical slice maps Hermes's existing API-request hooks to one Relay -LLM lifecycle per logical model call and aggregates terminal primary calls into -`hermes.model_call.count`. Retries sharing an API request ID remain one logical -call. The counter uses bounded provider family, model family, locality, and -outcome dimensions; raw model IDs and request or response content are never -included in the metrics event. - -The subscriber stores cumulative counters and a random opaque install ID in -`$HERMES_HOME/telemetry/shared_metrics/metrics.sqlite3`. After Relay's flush -barrier, Hermes commits un-packaged deltas to an immutable outbox record and -atomically writes a `hermes.shared_metrics.v1` JSON package under -`$HERMES_HOME/telemetry/shared_metrics/outbox/`. Repeating export without new -model calls reuses pending package IDs and does not package a count twice. The -proof mode does not perform network I/O. Task, tool, approval, and skill metrics -remain follow-on slices. The closed package contract ships with the plugin at -`schemas/hermes.shared_metrics.v1.schema.json`. - -When shared metrics are enabled without a separate rich-observability -configuration, the plugin does not emit its existing content-bearing turn, -model, tool, approval, or subagent events. ATOF, ATIF, adaptive components, and -dynamic plugins remain separate explicit configuration choices. - -### End-to-End Smoke Test - -The repository includes a deterministic smoke runner that exercises the real -Hermes CLI and native NeMo Relay binding without external model credentials. It -starts a loopback OpenAI-compatible model server, creates an isolated -`HERMES_HOME`, runs one `hermes chat` turn, and validates the resulting SQLite -counter and schema-conformant JSON package: - -```bash -./.venv/bin/python scripts/smoke_nemo_relay_shared_metrics.py \ - --relay-python ../nemo-relay/python -``` - -The NeMo Relay Python binding must already be built. By default, the runner -expects the Hermes checkout as its working directory and looks for a sibling -`nemo-relay` checkout. Use `--hermes-repo`, `--relay-python`, or `--output-dir` -to override those locations. The generated profile, captured CLI output, -SQLite database, and immutable package are retained in the artifact directory -printed after a successful run. +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 @@ -315,13 +250,10 @@ For the full generic Hermes middleware contract, see ## Canonical Local Examples -The observe-only examples in this section use a supported NeMo Relay 0.5.x -distribution and a local Ollama model served through the -OpenAI-compatible API. +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 -pip install "nemo-relay>=0.5.0,<0.6.0" - export HERMES_HOME=/tmp/hermes-nemo-relay-docs/hermes-home mkdir -p "$HERMES_HOME" diff --git a/plugins/observability/nemo_relay/__init__.py b/plugins/observability/nemo_relay/__init__.py index 895320b18830..d539e0eb198a 100644 --- a/plugins/observability/nemo_relay/__init__.py +++ b/plugins/observability/nemo_relay/__init__.py @@ -4,7 +4,6 @@ from __future__ import annotations import atexit import asyncio -import contextvars import inspect import json import logging @@ -13,22 +12,10 @@ import threading import tomllib from collections.abc import Callable from dataclasses import dataclass, field -from functools import wraps from pathlib import Path from typing import Any, Optional -from .shared_metrics import SharedMetricsStore -from .shared_metrics_contract import ( - MODEL_CALL_SCOPE as _SHARED_METRICS_MODEL_CALL_SCOPE, - SCHEMA_KEY as _SHARED_METRICS_SCHEMA_KEY, - SCHEMA_VERSION as _SHARED_METRICS_SCHEMA_VERSION, - SESSION_SCOPE as _SHARED_METRICS_SESSION_SCOPE, - SUBSCRIBER_NAME as _SHARED_METRICS_SUBSCRIBER_NAME, - execution_surface as _execution_surface, - model_call_fields as _model_call_fields, - model_call_outcome as _model_call_outcome, -) -from .shared_metrics_subscriber import SharedMetricsSubscriber +from hermes_cli.observability import relay_runtime logger = logging.getLogger(__name__) @@ -45,39 +32,24 @@ _RELAY_LLM_SURFACE_BY_API_MODE = { @dataclass class _SessionState: session_id: str - lock: threading.RLock = field(default_factory=threading.RLock, repr=False) - closing: bool = False + relay_session: relay_runtime.RelaySession | None = None handle: Any = None - metrics_handle: Any = None - metrics_context: contextvars.Context | None = None atif_exporter: Any = None atif_subscriber_name: str = "" is_embedded_subagent: bool = False parent_session_id: str = "" llm_spans: dict[str, Any] = field(default_factory=dict) tool_spans: dict[str, Any] = field(default_factory=dict) - metrics_model_calls: dict[str, "_SharedMetricsModelCall"] = field( - default_factory=dict - ) @dataclass -class _SharedMetricsModelCall: - handle: Any - task_id: str - fields: dict[str, str] - - -@dataclass -class _SubagentParent: +class _SubagentContext: parent_session_id: str - parent_handle: Any metadata: dict[str, Any] @dataclass class _Settings: - shared_metrics_enabled: bool = False plugins_toml_path: str = "" plugins_config: dict[str, Any] | None = None dynamic_plugins: list[dict[str, Any]] = field(default_factory=list) @@ -96,125 +68,27 @@ class _Settings: atif_model_name: str = "unknown" -def _with_session_state_lock(method: Callable[..., Any]) -> Callable[..., Any]: - """Serialize one session without blocking unrelated sessions.""" - - @wraps(method) - def wrapped( - self: "_Runtime", - event: dict[str, Any], - *args: Any, - **kwargs: Any, - ) -> Any: - with self._state_lock: - state = self.sessions.get(_session_id(event)) - if state is None: - return None - with state.lock: - if state.closing: - return None - return method(self, event, state, *args, **kwargs) - - return wrapped - - -def _with_ensured_session_lock(method: Callable[..., Any]) -> Callable[..., Any]: - """Create a missing session, then serialize work scoped to it.""" - - @wraps(method) - def wrapped( - self: "_Runtime", - event: dict[str, Any], - *args: Any, - **kwargs: Any, - ) -> Any: - state = self.ensure_session(event) - with state.lock: - if state.closing: - return None - return method(self, event, state, *args, **kwargs) - - return wrapped - - class _Runtime: - def __init__(self, nemo_relay: Any, settings: _Settings) -> None: + def __init__( + self, + nemo_relay: Any, + settings: _Settings, + host: relay_runtime.RelayRuntime, + ) -> None: self.nemo_relay = nemo_relay self.settings = settings - self._state_lock = threading.RLock() - self._plugin_lifecycle_lock = threading.RLock() + self.host = host self.sessions: dict[str, _SessionState] = {} - self.subagent_parents: dict[str, _SubagentParent] = {} + self.subagent_contexts: dict[str, _SubagentContext] = {} self.atof_exporter: Any = None self._atof_subscriber_name = "hermes.nemo_relay.atof" self._plugin_activation: Any = None self._shutdown_registered = False - self.shared_metrics: SharedMetricsSubscriber | None = None - if settings.shared_metrics_enabled: - try: - from hermes_cli import __version__ - - self.shared_metrics = SharedMetricsSubscriber( - SharedMetricsStore(), - __version__, - ) - except Exception: - logger.warning( - "NeMo Relay shared metrics disabled: local store initialization failed", - exc_info=True, - ) - self._shared_metrics_registered = False - self._configure_shared_metrics() self._plugin_config_initialized = self._configure_plugins_toml() self._plugin_config_needs_reinit = False if not self._plugin_config_initialized: self._activate_direct_fallbacks() - def _configure_shared_metrics(self) -> None: - if self.shared_metrics is None: - return - subscribers = getattr(self.nemo_relay, "subscribers", None) - register = getattr(subscribers, "register", None) - if not callable(register): - logger.warning( - "NeMo Relay shared metrics disabled: subscriber registration is unavailable" - ) - self.shared_metrics = None - return - try: - register(_SHARED_METRICS_SUBSCRIBER_NAME, self.shared_metrics) - except Exception as exc: - logger.warning( - "NeMo Relay shared metrics subscriber registration failed: %s", exc - ) - self.shared_metrics = None - return - self._shared_metrics_registered = True - self._ensure_shutdown_registered() - - def export_shared_metrics(self) -> list[Path]: - """Commit and export model-call deltas after Relay subscriber flush.""" - if self.shared_metrics is None: - return [] - try: - return self.shared_metrics.store.create_and_export_package() - except Exception: - logger.warning("Hermes shared-metrics package export failed", exc_info=True) - return [] - - def rich_observability_enabled(self) -> bool: - if not self.settings.shared_metrics_enabled: - return True - return bool( - self.settings.atof_enabled - or self.settings.atif_enabled - or _enabled_component_config( - self.settings.plugins_config, - "observability", - ) - is not None - ) - def _configure_plugins_toml(self) -> bool: if not self.settings.plugins_config: return False @@ -284,9 +158,7 @@ class _Runtime: # before its awaitable resolves, including error results. self._plugin_activation = None self._plugin_config_initialized = False - self._plugin_config_needs_reinit = bool( - self.settings.plugins_config - ) + self._plugin_config_needs_reinit = bool(self.settings.plugins_config) else: failures.append("dynamic plugin activation has no close method") else: @@ -370,268 +242,82 @@ class _Runtime: self.atof_exporter = None def ensure_session(self, kwargs: dict[str, Any]) -> _SessionState: - with self._plugin_lifecycle_lock: - self._maybe_reinitialize_plugins_toml() - with self._state_lock: - session_id = _session_id(kwargs) - state = self.sessions.get(session_id) - if state is not None: - self._ensure_shared_metrics_session(state, kwargs) - return state + self._maybe_reinitialize_plugins_toml() + session_id = _session_id(kwargs) + state = self.sessions.get(session_id) + if state is not None: + return state - state = _SessionState(session_id=session_id) - if self.rich_observability_enabled(): - 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.{session_id}" - ) - state.atif_exporter.register(state.atif_subscriber_name) - - subagent_parent = self.subagent_parents.get(session_id) - metadata = _metadata(kwargs) - parent_handle = None - if subagent_parent is not None: - parent_handle = subagent_parent.parent_handle - metadata = {**metadata, **subagent_parent.metadata} - state.is_embedded_subagent = True - state.parent_session_id = subagent_parent.parent_session_id - - state.handle = self.nemo_relay.scope.push( - f"hermes-session-{session_id}", - self.nemo_relay.ScopeType.Agent, - handle=parent_handle, - data={"session_id": session_id}, - metadata=metadata, - ) - - self._ensure_shared_metrics_session(state, kwargs) - self.sessions[session_id] = state - return state - - def _ensure_shared_metrics_session( - self, - state: _SessionState, - kwargs: dict[str, Any], - ) -> None: - if ( - state.closing - or not self._shared_metrics_registered - or state.metrics_handle is not None - ): - return - metrics_context = contextvars.Context() - try: - state.metrics_handle = metrics_context.run( - self.nemo_relay.scope.push, - _SHARED_METRICS_SESSION_SCOPE, - self.nemo_relay.ScopeType.Agent, - input={"execution_surface": _execution_surface(kwargs)}, - metadata={_SHARED_METRICS_SCHEMA_KEY: _SHARED_METRICS_SCHEMA_VERSION}, + 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"}, ) - except Exception: - logger.warning( - "NeMo Relay shared-metrics session start failed", exc_info=True - ) - return - state.metrics_context = metrics_context + state.atif_subscriber_name = f"hermes.nemo_relay.atif.{session_id}" + state.atif_exporter.register(state.atif_subscriber_name) - def _run_in_metrics_context( + rich_metadata = _metadata(kwargs) + subagent_context = self.subagent_contexts.get(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": 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 + self.sessions[session_id] = state + return state + + def run_in_session( self, state: _SessionState, callback: Callable[..., Any], *args: Any, **kwargs: Any, ) -> Any: - """Run a native lifecycle call against the isolated metrics stack.""" - if state.metrics_context is None: - raise RuntimeError("shared-metrics scope context is unavailable") - - def invoke() -> Any: - # Relay's manual LLM helpers call the native API directly. Re-sync - # this Context's stack before entering them so scope-local policy - # cannot leak in from the rich-observability stack on this thread. - self.nemo_relay.get_scope_stack() - return callback(*args, **kwargs) - - return state.metrics_context.run(invoke) - - @_with_ensured_session_lock - def start_model_call(self, kwargs: dict[str, Any], state: _SessionState) -> None: - if not self._shared_metrics_registered or state.metrics_handle is None: - return - request_id = str(kwargs.get("api_request_id") or "") - if not request_id: - return - fields = _model_call_fields(kwargs) - model_family = fields.pop("model_family") - existing = state.metrics_model_calls.get(request_id) - if existing is not None: - # A logical call can span retries or provider fallback. Attribute its - # terminal provider fields to the most recent attempt without opening - # a second Relay lifecycle. Relay's model_name remains the logical - # model family recorded when the lifecycle started. - existing.fields = fields - return - request = self.nemo_relay.LLMRequest({}, {}) - try: - handle = self._run_in_metrics_context( - state, - self.nemo_relay.llm.call, - _SHARED_METRICS_MODEL_CALL_SCOPE, - request, - handle=state.metrics_handle, - metadata={_SHARED_METRICS_SCHEMA_KEY: _SHARED_METRICS_SCHEMA_VERSION}, - model_name=model_family, - ) - except Exception: - logger.warning( - "NeMo Relay shared-metrics model-call start failed", exc_info=True - ) - return - state.metrics_model_calls[request_id] = _SharedMetricsModelCall( - handle=handle, - task_id=str(kwargs.get("task_id") or ""), - fields=fields, + 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 _finish_model_call( - self, - state: _SessionState, - request_id: str, - outcome: str, - ) -> None: - model_call = state.metrics_model_calls.pop(request_id, None) - if model_call is None: - return - try: - self._run_in_metrics_context( - state, - self.nemo_relay.llm.call_end, - model_call.handle, - {**model_call.fields, "outcome": outcome}, - metadata={_SHARED_METRICS_SCHEMA_KEY: _SHARED_METRICS_SCHEMA_VERSION}, - ) - except Exception: - logger.warning( - "NeMo Relay shared-metrics model-call end failed", exc_info=True - ) - - @_with_session_state_lock - def end_model_call(self, kwargs: dict[str, Any], state: _SessionState) -> None: - request_id = str(kwargs.get("api_request_id") or "") - model_call = state.metrics_model_calls.get(request_id) - if model_call is None: - return - fields = _model_call_fields(kwargs) - fields.pop("model_family") - model_call.fields = fields - self._finish_model_call(state, request_id, _model_call_outcome(kwargs)) - - @_with_session_state_lock - def end_pending_model_calls( - self, - kwargs: dict[str, Any], - state: _SessionState, - ) -> None: - self._end_pending_model_calls(state, kwargs) - - def _end_pending_model_calls( - self, - state: _SessionState, - kwargs: dict[str, Any], - ) -> None: - task_id = str(kwargs.get("task_id") or "") - request_ids = [ - request_id - for request_id, model_call in state.metrics_model_calls.items() - if not task_id or model_call.task_id == task_id - ] - outcome = "cancelled" if kwargs.get("interrupted") else "failed" - for request_id in request_ids: - self._finish_model_call(state, request_id, outcome) - 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" - ): + 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 _clear_static_plugins_after_last_session(self) -> None: - """Clear static plugin state without holding the session-map lock.""" - with self._plugin_lifecycle_lock: - with self._state_lock: - if self.sessions or self._plugin_activation is not None: - return - should_clear = self._plugin_config_initialized - if not should_clear and self.settings.plugins_config: - self._plugin_config_needs_reinit = True - if should_clear: - self._clear_plugins_toml() + 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]) -> None: session_id = _session_id(kwargs) - with self._state_lock: - self.subagent_parents.pop(session_id, None) - state = self.sessions.get(session_id) + self.subagent_contexts.pop(session_id, None) + state = self.sessions.pop(session_id, None) if state is None: return failures: list[str] = [] - with state.lock: - if state.closing: - return - state.closing = True - self._end_pending_model_calls(state, {"session_id": session_id}) - if state.metrics_handle is not None and state.metrics_context is not None: - try: - self._run_in_metrics_context( - state, - self.nemo_relay.scope.pop, - state.metrics_handle, - output={}, - metadata={ - _SHARED_METRICS_SCHEMA_KEY: _SHARED_METRICS_SCHEMA_VERSION - }, - ) - except Exception as exc: - failures.append(f"shared-metrics session scope pop failed: {exc}") - if state.handle is not None: - try: - self.nemo_relay.scope.pop(state.handle, output=_jsonable(kwargs)) - except Exception as exc: - failures.append(f"session scope pop failed: {exc}") try: - _flush_relay_subscribers(self.nemo_relay) + self.host.close_session(kwargs) except Exception as exc: - failures.append(f"subscriber flush failed: {exc}") - self.export_shared_metrics() + failures.append(f"core session close failed: {exc}") try: self.export_atif(state) except Exception as exc: @@ -641,13 +327,21 @@ class _Runtime: state.atif_exporter.deregister(state.atif_subscriber_name) except Exception as exc: failures.append(f"ATIF deregister failed: {exc}") - with self._state_lock: - if self.sessions.get(session_id) is state: - self.sessions.pop(session_id, None) - try: - self._clear_static_plugins_after_last_session() - except Exception as exc: - failures.append(f"plugin configuration clear failed: {exc}") + 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", @@ -658,38 +352,17 @@ class _Runtime: def shutdown(self) -> None: """Close active sessions and the process-lifetime plugin activation.""" failures: list[str] = [] - with self._state_lock: - session_ids = list(self.sessions) - for session_id in session_ids: + for session_id in list(self.sessions): try: - self.close_session({ - "session_id": session_id, - "reason": "runtime_shutdown", - }) + self.close_session({"session_id": session_id, "reason": "runtime_shutdown"}) except Exception as exc: failures.append(f"session {session_id} close failed: {exc}") - with self._plugin_lifecycle_lock: - if self._plugin_config_initialized: - try: - self._clear_plugins_toml() - except Exception as exc: - failures.append(f"plugin runtime 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._shared_metrics_registered: - try: - _flush_relay_subscribers(self.nemo_relay) - except Exception as exc: - failures.append(f"shared-metrics subscriber flush failed: {exc}") - self.export_shared_metrics() - try: - subscribers = getattr(self.nemo_relay, "subscribers", None) - deregister = getattr(subscribers, "deregister", None) - if callable(deregister): - deregister(_SHARED_METRICS_SUBSCRIBER_NAME) - except Exception as exc: - failures.append(f"shared-metrics subscriber deregister failed: {exc}") - finally: - self._shared_metrics_registered = False if self._shutdown_registered and self._plugin_activation is None: atexit.unregister(self.shutdown) self._shutdown_registered = False @@ -701,7 +374,9 @@ class _Runtime: def mark(self, name: str, kwargs: dict[str, Any]) -> None: state = self.ensure_session(kwargs) - self.nemo_relay.scope.event( + self.run_in_session( + state, + self.nemo_relay.scope.event, name, handle=state.handle, data=_jsonable(kwargs), @@ -709,16 +384,18 @@ class _Runtime: ) def mark_subagent_start(self, kwargs: dict[str, Any]) -> None: + self.host.register_subagent(kwargs) parent_state = self.ensure_session(kwargs) metadata = _metadata(kwargs) child_session_id = _child_session_id(kwargs) if child_session_id: - self.subagent_parents[child_session_id] = _SubagentParent( + self.subagent_contexts[child_session_id] = _SubagentContext( parent_session_id=parent_state.session_id, - parent_handle=parent_state.handle, metadata=_subagent_child_metadata(kwargs, metadata), ) - self.nemo_relay.scope.event( + self.run_in_session( + parent_state, + self.nemo_relay.scope.event, "hermes.subagent.start", handle=parent_state.handle, data=_jsonable(kwargs), @@ -726,25 +403,23 @@ class _Runtime: ) def mark_subagent_stop(self, kwargs: dict[str, Any]) -> None: + self.host.unregister_subagent(kwargs) child_session_id = _child_session_id(kwargs) if child_session_id: - self.subagent_parents.pop(child_session_id, None) + self.subagent_contexts.pop(child_session_id, None) self.mark("hermes.subagent.stop", kwargs) def managed_llm_enabled(self) -> bool: return ( (self.settings.adaptive_enabled or self._plugin_activation is not None) - and callable( - getattr(getattr(self.nemo_relay, "llm", None), "execute", None) - ) + and callable(getattr(getattr(self.nemo_relay, "llm", None), "execute", None)) and callable(getattr(self.nemo_relay, "LLMRequest", None)) ) def managed_tool_enabled(self) -> bool: return ( - self.settings.adaptive_enabled or self._plugin_activation is not None - ) and callable( - getattr(getattr(self.nemo_relay, "tools", None), "execute", None) + (self.settings.adaptive_enabled or self._plugin_activation is not None) + and callable(getattr(getattr(self.nemo_relay, "tools", None), "execute", None)) ) def _run_managed_with_downstream_preservation( @@ -782,9 +457,7 @@ class _Runtime: try: managed_result = _resolve_awaitable(make_managed_execute(_impl)) except Exception as exc: - if downstream_error is not None and _is_relay_wrapped_callback_error( - exc, callback_error - ): + if downstream_error is not None and _is_relay_wrapped_callback_error(exc, callback_error): raise downstream_error raise if ( @@ -809,17 +482,21 @@ class _Runtime: def _make_managed(impl: Callable[[Any], Any]) -> Any: async def _managed_execute() -> Any: - result = self.nemo_relay.llm.execute( + result = self.run_in_session( + state, + self.nemo_relay.llm.execute, _relay_llm_surface(kwargs), request, impl, handle=state.handle, - data=_jsonable({ - "turn_id": kwargs.get("turn_id"), - "api_request_id": kwargs.get("api_request_id"), - "api_call_count": kwargs.get("api_call_count"), - "mode": self.settings.adaptive_mode, - }), + data=_jsonable( + { + "turn_id": kwargs.get("turn_id"), + "api_request_id": kwargs.get("api_request_id"), + "api_call_count": kwargs.get("api_call_count"), + "mode": self.settings.adaptive_mode, + } + ), metadata=_metadata(kwargs), model_name=str(kwargs.get("model") or ""), ) @@ -830,11 +507,7 @@ class _Runtime: return _managed_execute() return self._run_managed_with_downstream_preservation( - next_call, - _normalize, - _llm_response_payload, - _make_managed, - preserve_raw_response=True, + next_call, _normalize, _llm_response_payload, _make_managed, preserve_raw_response=True ) def execute_tool(self, kwargs: dict[str, Any]) -> Any: @@ -850,17 +523,21 @@ class _Runtime: def _make_managed(impl: Callable[[Any], Any]) -> Any: async def _managed_execute() -> Any: - result = self.nemo_relay.tools.execute( + result = self.run_in_session( + state, + self.nemo_relay.tools.execute, tool_name, args, impl, handle=state.handle, - data=_jsonable({ - "turn_id": kwargs.get("turn_id"), - "api_request_id": kwargs.get("api_request_id"), - "tool_call_id": kwargs.get("tool_call_id"), - "mode": self.settings.adaptive_mode, - }), + data=_jsonable( + { + "turn_id": kwargs.get("turn_id"), + "api_request_id": kwargs.get("api_request_id"), + "tool_call_id": kwargs.get("tool_call_id"), + "mode": self.settings.adaptive_mode, + } + ), metadata=_metadata(kwargs), ) if inspect.isawaitable(result): @@ -906,16 +583,8 @@ def on_session_start(**kwargs: Any) -> None: def on_session_end(**kwargs: Any) -> None: runtime = _get_runtime() - if runtime is None: - return - _safe(lambda: runtime.end_pending_model_calls(kwargs)) - if runtime.rich_observability_enabled(): - _safe( - lambda: ( - runtime.mark("hermes.session.end", kwargs), - runtime.export_atif(runtime.ensure_session(kwargs)), - ) - ) + 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: @@ -932,13 +601,13 @@ def on_session_reset(**kwargs: Any) -> None: def on_pre_llm_call(**kwargs: Any) -> None: runtime = _get_runtime() - if runtime is not None and runtime.rich_observability_enabled(): + 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 and runtime.rich_observability_enabled(): + if runtime is not None: _safe(lambda: runtime.mark("hermes.turn.end", kwargs)) @@ -946,27 +615,21 @@ def on_pre_api_request(**kwargs: Any) -> None: runtime = _get_runtime() if runtime is None: return - _safe(lambda: runtime.start_model_call(kwargs)) - if not runtime.rich_observability_enabled(): - return if runtime.managed_llm_enabled(): return def _record() -> None: state = runtime.ensure_session(kwargs) request_payload = kwargs.get("request") - request_body = ( - request_payload.get("body") if isinstance(request_payload, dict) else {} - ) + request_body = request_payload.get("body") if isinstance(request_payload, dict) else {} request = runtime.nemo_relay.LLMRequest({}, _jsonable(request_body)) - span = runtime.nemo_relay.llm.call( + span = runtime.run_in_session( + state, + runtime.nemo_relay.llm.call, str(kwargs.get("provider") or "llm"), request, handle=state.handle, - data=_jsonable({ - "turn_id": kwargs.get("turn_id"), - "api_request_id": kwargs.get("api_request_id"), - }), + data=_jsonable({"turn_id": kwargs.get("turn_id"), "api_request_id": kwargs.get("api_request_id")}), metadata=_metadata(kwargs), model_name=str(kwargs.get("model") or ""), ) @@ -979,9 +642,6 @@ def on_post_api_request(**kwargs: Any) -> None: runtime = _get_runtime() if runtime is None: return - _safe(lambda: runtime.end_model_call({**kwargs, "outcome": "success"})) - if not runtime.rich_observability_enabled(): - return if runtime.managed_llm_enabled(): return @@ -991,13 +651,12 @@ def on_post_api_request(**kwargs: Any) -> None: if span is None: runtime.mark("hermes.api.response.unmatched", kwargs) return - runtime.nemo_relay.llm.call_end( + runtime.run_in_session( + state, + runtime.nemo_relay.llm.call_end, span, _jsonable(kwargs.get("response") or {}), - data=_jsonable({ - "usage": kwargs.get("usage"), - "finish_reason": kwargs.get("finish_reason"), - }), + data=_jsonable({"usage": kwargs.get("usage"), "finish_reason": kwargs.get("finish_reason")}), metadata=_metadata(kwargs), ) @@ -1006,7 +665,7 @@ def on_post_api_request(**kwargs: Any) -> None: def on_api_request_error(**kwargs: Any) -> None: runtime = _get_runtime() - if runtime is None or not runtime.rich_observability_enabled(): + if runtime is None: return if runtime.managed_llm_enabled(): return @@ -1017,7 +676,9 @@ def on_api_request_error(**kwargs: Any) -> None: if span is None: runtime.mark("hermes.api.error", kwargs) return - runtime.nemo_relay.llm.call_end( + runtime.run_in_session( + state, + runtime.nemo_relay.llm.call_end, span, {"error": _jsonable(kwargs.get("error") or {})}, data=_jsonable(kwargs), @@ -1031,21 +692,18 @@ def on_pre_tool_call(**kwargs: Any) -> None: runtime = _get_runtime() if runtime is None: return - if not runtime.rich_observability_enabled(): - return if runtime.managed_tool_enabled(): return def _record() -> None: state = runtime.ensure_session(kwargs) - span = runtime.nemo_relay.tools.call( + span = runtime.run_in_session( + state, + runtime.nemo_relay.tools.call, str(kwargs.get("tool_name") or "tool"), _jsonable(kwargs.get("args") or {}), handle=state.handle, - data=_jsonable({ - "turn_id": kwargs.get("turn_id"), - "api_request_id": kwargs.get("api_request_id"), - }), + data=_jsonable({"turn_id": kwargs.get("turn_id"), "api_request_id": kwargs.get("api_request_id")}), metadata=_metadata(kwargs), tool_call_id=str(kwargs.get("tool_call_id") or ""), ) @@ -1058,8 +716,6 @@ def on_post_tool_call(**kwargs: Any) -> None: runtime = _get_runtime() if runtime is None: return - if not runtime.rich_observability_enabled(): - return if runtime.managed_tool_enabled(): return @@ -1069,13 +725,12 @@ def on_post_tool_call(**kwargs: Any) -> None: if span is None: runtime.mark("hermes.tool.response.unmatched", kwargs) return - runtime.nemo_relay.tools.call_end( + runtime.run_in_session( + state, + runtime.nemo_relay.tools.call_end, span, _jsonable(kwargs.get("result")), - data=_jsonable({ - "status": kwargs.get("status"), - "duration_ms": kwargs.get("duration_ms"), - }), + data=_jsonable({"status": kwargs.get("status"), "duration_ms": kwargs.get("duration_ms")}), metadata=_metadata(kwargs), ) @@ -1084,29 +739,25 @@ def on_post_tool_call(**kwargs: Any) -> None: def on_pre_approval_request(**kwargs: Any) -> None: runtime = _get_runtime() - if runtime is None: - return - if runtime.rich_observability_enabled(): + 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 None: - return - if runtime.rich_observability_enabled(): + 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 and runtime.rich_observability_enabled(): + 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 and runtime.rich_observability_enabled(): + if runtime is not None: _safe(lambda: runtime.mark_subagent_stop(kwargs)) @@ -1140,17 +791,16 @@ def _get_runtime() -> Optional[_Runtime]: if isinstance(_RUNTIME, _Runtime): return _RUNTIME try: - import nemo_relay as nemo_runtime - except Exception as exc: - logger.debug("NeMo Relay plugin disabled: import failed: %s", exc) - _RUNTIME = _INIT_FAILED - return None - try: - _RUNTIME = _Runtime(nemo_relay=nemo_runtime, settings=_load_settings()) - except Exception as exc: - logger.debug( - "NeMo Relay plugin disabled: init failed: %s", exc, exc_info=True + host = relay_runtime.get_runtime() + if host is None: + raise RuntimeError("Hermes core Relay runtime is unavailable") + _RUNTIME = _Runtime( + nemo_relay=host.relay, + settings=_load_settings(), + host=host, ) + except Exception as exc: + logger.debug("NeMo Relay plugin disabled: init failed: %s", exc, exc_info=True) _RUNTIME = _INIT_FAILED return None return _RUNTIME @@ -1161,7 +811,6 @@ def _load_settings() -> _Settings: plugins_config = _load_plugins_config(plugins_toml_path) adaptive_config = _enabled_component_config(plugins_config, "adaptive") return _Settings( - shared_metrics_enabled=_shared_metrics_enabled(), plugins_toml_path=plugins_toml_path, plugins_config=plugins_config, dynamic_plugins=_dynamic_plugin_specs(plugins_config, plugins_toml_path), @@ -1173,8 +822,7 @@ def _load_settings() -> _Settings: 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_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", @@ -1182,31 +830,6 @@ def _load_settings() -> _Settings: ) -def _shared_metrics_enabled() -> bool: - try: - from hermes_cli.config import load_config_readonly - - config = load_config_readonly() or {} - except Exception: - logger.debug( - "Unable to read Hermes shared-metrics configuration", exc_info=True - ) - return False - if not isinstance(config, dict): - return False - plugins = config.get("plugins") - if not isinstance(plugins, dict): - return False - entries = plugins.get("entries") - if not isinstance(entries, dict): - return False - entry = entries.get("observability/nemo_relay") or entries.get("nemo_relay") - if not isinstance(entry, dict): - return False - shared_metrics = entry.get("shared_metrics") - return isinstance(shared_metrics, dict) and shared_metrics.get("enabled") is True - - def _static_plugin_config(plugins_config: dict[str, Any]) -> dict[str, Any]: """Return Relay's base config without embedding- or gateway-host fields.""" return { @@ -1279,15 +902,13 @@ def _dynamic_plugin_specs( 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 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 NeMo Relay dynamic_plugins[%d]: config must be an object", index ) invalid = True continue @@ -1304,9 +925,7 @@ def _dynamic_plugin_specs( spec: dict[str, Any] = { "plugin_id": plugin_id.strip(), "kind": kind, - "manifest_ref": _config_relative_path( - manifest_ref.strip(), plugins_toml_path - ), + "manifest_ref": _config_relative_path(manifest_ref.strip(), plugins_toml_path), "config": config, } if environment_ref is not None: @@ -1328,9 +947,7 @@ def _config_relative_path(value: str, plugins_toml_path: str) -> str: 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" - ) + 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) @@ -1421,7 +1038,8 @@ def _child_session_id(kwargs: dict[str, Any]) -> str: def _subagent_child_metadata( - kwargs: dict[str, Any], parent_metadata: dict[str, Any] + kwargs: dict[str, Any], + parent_metadata: dict[str, Any], ) -> dict[str, Any]: child_session_id = _child_session_id(kwargs) metadata = { @@ -1447,10 +1065,7 @@ def _subagent_child_metadata( def _api_key(kwargs: dict[str, Any]) -> str: - return str( - kwargs.get("api_request_id") - or f"{_session_id(kwargs)}:{kwargs.get('api_call_count') or 'api'}" - ) + return str(kwargs.get("api_request_id") or f"{_session_id(kwargs)}:{kwargs.get('api_call_count') or 'api'}") def _tool_key(kwargs: dict[str, Any]) -> str: @@ -1530,10 +1145,13 @@ def _jsonable(value: Any) -> Any: def _json_semantically_equal(left: Any, right: Any) -> bool: """Compare JSON-compatible values without conflating booleans and numbers.""" try: - options = {"ensure_ascii": False, "sort_keys": True, "separators": (",", ":")} - return json.dumps(_jsonable(left), **options) == json.dumps( - _jsonable(right), **options + left_json = json.dumps( + _jsonable(left), ensure_ascii=False, sort_keys=True, separators=(",", ":") ) + right_json = json.dumps( + _jsonable(right), ensure_ascii=False, sort_keys=True, separators=(",", ":") + ) + return left_json == right_json except (TypeError, ValueError): return False @@ -1548,16 +1166,12 @@ def _original_downstream_error(exc: Exception) -> BaseException: # Hermes wraps downstream execution failures in a local/private exception # class, so detect the wrapper by shape instead of importing it here. original = getattr(exc, "original", None) - if exc.__class__.__name__ == "_DownstreamExecutionError" and isinstance( - original, BaseException - ): + if exc.__class__.__name__ == "_DownstreamExecutionError" and isinstance(original, BaseException): return original return exc -def _is_relay_wrapped_callback_error( - exc: Exception, callback_error: Exception | None -) -> bool: +def _is_relay_wrapped_callback_error(exc: Exception, callback_error: Exception | None) -> bool: # NeMo Relay re-wraps a failing callback as ``RuntimeError("internal error: # : ")``. Match by prefix rather than exact equality so a # trailing traceback/suffix in a future Relay version doesn't silently defeat @@ -1598,25 +1212,13 @@ def _llm_response_payload(response: Any) -> Any: if reasoning is not None: assistant_message["reasoning_content"] = _jsonable(reasoning) elif isinstance(payload, dict): - assistant_message["content"] = ( - payload.get("content") or payload.get("output_text") or "" - ) + assistant_message["content"] = payload.get("content") or payload.get("output_text") or "" return { - "model": _value( - response, - "model", - payload.get("model") if isinstance(payload, dict) else None, - ), + "model": _value(response, "model", payload.get("model") if isinstance(payload, dict) else None), "assistant_message": assistant_message, "finish_reason": finish_reason, - "usage": _jsonable( - _value( - response, - "usage", - payload.get("usage") if isinstance(payload, dict) else None, - ) - ), + "usage": _jsonable(_value(response, "usage", payload.get("usage") if isinstance(payload, dict) else None)), } @@ -1626,14 +1228,16 @@ def _tool_calls_payload(tool_calls: Any) -> list[dict[str, Any]]: normalized: list[dict[str, Any]] = [] for call in tool_calls: function = _value(call, "function") - normalized.append({ - "id": _value(call, "id"), - "type": _value(call, "type", "function") or "function", - "function": { - "name": _value(function, "name"), - "arguments": _value(function, "arguments"), - }, - }) + normalized.append( + { + "id": _value(call, "id"), + "type": _value(call, "type", "function") or "function", + "function": { + "name": _value(function, "name"), + "arguments": _value(function, "arguments"), + }, + } + ) return normalized diff --git a/pyproject.toml b/pyproject.toml index 19950019f4c0..5031122d16cd 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -138,6 +138,10 @@ dependencies = [ # Hence the ``sys_platform == 'win32'`` marker: the dep (and its portalocker # / pywin32 tree) ships only where it's actually used. "concurrent-log-handler==0.9.29; sys_platform == 'win32'", + # First-party lifecycle and shared-metrics runtime. Relay 0.5 publishes + # native wheels for these targets only; keep unsupported Hermes targets + # installable without moving Relay back behind a user-selected extra. + "nemo-relay==0.5.0; (sys_platform == 'darwin' and platform_machine == 'arm64') or (sys_platform == 'linux' and platform_machine == 'x86_64') or (sys_platform == 'linux' and platform_machine == 'aarch64') or (sys_platform == 'win32' and platform_machine == 'AMD64') or (sys_platform == 'win32' and platform_machine == 'ARM64')", ] [project.optional-dependencies] @@ -205,7 +209,6 @@ vision = [] # extra that exposes a Starlette-backed server surface so pip/uv can't resolve # a vulnerable pre-1.0.1 transitive. Bump in lockstep with uv.lock. mcp = ["mcp==1.26.0", "starlette==1.0.1"] # starlette: CVE-2026-48710 -nemo-relay = ["nemo-relay>=0.5.0,<0.6.0"] homeassistant = ["aiohttp==3.14.1"] sms = ["aiohttp==3.14.1"] teams = ["microsoft-teams-apps==2.0.13.4", "aiohttp==3.14.1"] # aiohttp 3.14.1: CVE-2026-34993(RCE)/47265 + 34513/34518/34519/34520/34525 @@ -336,7 +339,7 @@ locales = ["locales/*.yaml"] "optional-mcps/n8n" = ["optional-mcps/n8n/manifest.yaml"] [tool.setuptools.package-data] -hermes_cli = ["web_dist/**/*", "tui_dist/**/*", "scripts/install.sh", "scripts/install.ps1"] +hermes_cli = ["web_dist/**/*", "tui_dist/**/*", "scripts/install.sh", "scripts/install.ps1", "observability/schemas/*.json"] gateway = ["assets/**/*"] plugins = [ "*/dashboard/manifest.json", @@ -351,7 +354,6 @@ plugins = [ "**/plugin.yaml", "**/plugin.yml", "**/README.md", - "**/schemas/*.json", ] [tool.setuptools.packages.find] diff --git a/scripts/smoke_nemo_relay_shared_metrics.py b/scripts/smoke_nemo_relay_shared_metrics.py index 8ea4dde1b200..76e0a59538d4 100644 --- a/scripts/smoke_nemo_relay_shared_metrics.py +++ b/scripts/smoke_nemo_relay_shared_metrics.py @@ -152,7 +152,7 @@ def _arguments() -> argparse.Namespace: "--relay-python", type=Path, default=None, - help="NeMo Relay checkout's python directory", + help="Optional NeMo Relay checkout's python directory", ) parser.add_argument( "--output-dir", @@ -174,13 +174,9 @@ def _write_config(home: Path, port: int) -> None: api_key: no-key-required security: tirith_enabled: false -plugins: - enabled: - - observability/nemo_relay - entries: - observability/nemo_relay: - shared_metrics: - enabled: true +telemetry: + shared_metrics: + enabled: true """, encoding="utf-8", ) @@ -270,15 +266,13 @@ def _validate_package(outbox: Path, schema_path: Path) -> tuple[Path, dict[str, def main() -> int: args = _arguments() hermes_repo = args.hermes_repo.resolve() - relay_python = ( - args.relay_python.resolve() - if args.relay_python - else (hermes_repo.parent / "nemo-relay" / "python").resolve() - ) + relay_python = args.relay_python.resolve() if args.relay_python else None hermes = hermes_repo / ".venv" / "bin" / "hermes" if not hermes.is_file(): raise SystemExit(f"Hermes executable not found: {hermes}") - if not any((relay_python / "nemo_relay").glob("_native.*")): + if relay_python is not None and not any( + (relay_python / "nemo_relay").glob("_native.*") + ): raise SystemExit( "Built NeMo Relay Python binding not found under " f"{relay_python}; run the Relay Python build first" @@ -305,10 +299,11 @@ def main() -> int: _write_config(home, server.server_port) env = os.environ.copy() env["HERMES_HOME"] = str(home) - env["PYTHONPATH"] = os.pathsep.join([ - str(relay_python), - env.get("PYTHONPATH", ""), - ]).rstrip(os.pathsep) + if relay_python is not None: + env["PYTHONPATH"] = os.pathsep.join([ + str(relay_python), + env.get("PYTHONPATH", ""), + ]).rstrip(os.pathsep) result = subprocess.run( [ str(hermes), @@ -359,9 +354,8 @@ def main() -> int: package_path, package = _validate_package( telemetry / "outbox", hermes_repo - / "plugins" + / "hermes_cli" / "observability" - / "nemo_relay" / "schemas" / "hermes.shared_metrics.v1.schema.json", ) diff --git a/tests/plugins/test_nemo_relay_shared_metrics.py b/tests/hermes_cli/test_relay_shared_metrics.py similarity index 97% rename from tests/plugins/test_nemo_relay_shared_metrics.py rename to tests/hermes_cli/test_relay_shared_metrics.py index 5c2c1a87b814..a3bc18e8fbf3 100644 --- a/tests/plugins/test_nemo_relay_shared_metrics.py +++ b/tests/hermes_cli/test_relay_shared_metrics.py @@ -15,8 +15,8 @@ from types import SimpleNamespace from typing import Any import pytest -from plugins.observability.nemo_relay.shared_metrics import SharedMetricsStore -from plugins.observability.nemo_relay.shared_metrics_contract import ( +from hermes_cli.observability.shared_metrics import SharedMetricsStore +from hermes_cli.observability.shared_metrics_contract import ( MODEL_FAMILIES, MODEL_LOCALITIES, MODEL_OUTCOMES, @@ -33,9 +33,8 @@ from plugins.observability.nemo_relay.shared_metrics_contract import ( SCHEMA_PATH = ( Path(__file__).resolve().parents[2] - / "plugins" + / "hermes_cli" / "observability" - / "nemo_relay" / "schemas" / "hermes.shared_metrics.v1.schema.json" ) @@ -211,6 +210,12 @@ def test_model_family_accepts_only_allowlisted_declared_metadata(): assert model_family({"model": "private", "model_family": "private"}) == "unknown" +def test_model_family_prefers_the_provider_reported_terminal_model(): + assert ( + model_family({"model": "gpt-5", "response_model": "claude-sonnet"}) == "claude" + ) + + @pytest.mark.parametrize( ("platform", "expected"), [ @@ -245,6 +250,7 @@ def test_subscriber_contract_rejects_unknown_fields_and_dimension_values(): data={ "call_role": "primary", "locality": "remote", + "model_family": "gpt", "outcome": "success", "provider_family": "direct", }, diff --git a/tests/hermes_cli/test_relay_shared_metrics_runtime.py b/tests/hermes_cli/test_relay_shared_metrics_runtime.py new file mode 100644 index 000000000000..bbd612a4d5c6 --- /dev/null +++ b/tests/hermes_cli/test_relay_shared_metrics_runtime.py @@ -0,0 +1,463 @@ +"""Tests for the direct Hermes-to-Relay shared-metrics runtime.""" + +from __future__ import annotations + +import contextvars +import json +import threading +from pathlib import Path +from types import SimpleNamespace +from typing import Any + +import pytest + +from hermes_cli import plugins +from hermes_cli.observability import relay_runtime, relay_shared_metrics +from hermes_cli.plugins import PluginManager + + +class _Request: + def __init__(self, headers: dict[str, Any], content: dict[str, Any]) -> None: + self.headers = headers + self.content = content + + +class _Relay: + def __init__(self) -> None: + self.events: list[tuple[Any, ...]] = [] + self._callbacks: dict[str, Any] = {} + self._starts: dict[Any, dict[str, Any]] = {} + self._scope = contextvars.ContextVar("relay_scope", default=None) + self._scope_serial = 0 + self.ScopeType = SimpleNamespace(Agent="agent") + self.LLMRequest = _Request + 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) + self.subscribers = SimpleNamespace( + register=self._register, + deregister=self._deregister, + flush=self._flush, + ) + self.get_scope_stack = self._get_scope_stack + + def _scope_push(self, name: str, scope_type: Any, **kwargs: Any) -> Any: + self._scope_serial += 1 + handle = ("scope", name, self._scope_serial) + self._scope.set(handle) + self.events.append(("scope.push", name, scope_type, kwargs)) + return handle + + def _scope_pop(self, handle: Any, **kwargs: Any) -> None: + self.events.append(("scope.pop", handle, kwargs)) + + def _scope_event(self, name: str, **kwargs: Any) -> None: + self.events.append(("scope.event", name, kwargs)) + + def _get_scope_stack(self) -> Any: + current = self._scope.get() + self.events.append(("scope.sync", current)) + return current + + def _llm_call( + self, + name: str, + request: _Request, + **kwargs: Any, + ) -> Any: + handle = ("llm", name, len(self._starts)) + self._starts[handle] = kwargs + self.events.append(("llm.call", name, request.content, kwargs)) + return handle + + def _llm_call_end( + self, + handle: Any, + response: dict[str, Any], + **kwargs: Any, + ) -> None: + start = self._starts.pop(handle) + self.events.append(("llm.call_end", handle, response, kwargs)) + event = SimpleNamespace( + kind="scope", + category="llm", + name=handle[1], + scope_category="end", + category_profile={"model_name": start["model_name"]}, + metadata={ + **start["metadata"], + **kwargs["metadata"], + "otel.status_code": "OK", + }, + data=response, + ) + for callback in list(self._callbacks.values()): + callback(event) + + def _register(self, name: str, callback: Any) -> None: + self._callbacks[name] = callback + self.events.append(("subscribers.register", name)) + + def _deregister(self, name: str) -> None: + self._callbacks.pop(name, None) + self.events.append(("subscribers.deregister", name)) + + def _flush(self) -> None: + self.events.append(("subscribers.flush",)) + + +@pytest.fixture +def direct_runtime(tmp_path, monkeypatch): + fake = _Relay() + monkeypatch.setenv("HERMES_HOME", str(tmp_path / "hermes-home")) + monkeypatch.setattr(relay_runtime, "_load_nemo_relay", lambda: fake) + monkeypatch.setattr( + "hermes_cli.config.load_config_readonly", + lambda: {"telemetry": {"shared_metrics": {"enabled": True}}}, + ) + relay_shared_metrics._reset_for_tests() + relay_runtime._reset_for_tests() + monkeypatch.setattr(plugins, "_plugin_manager", PluginManager()) + yield fake + relay_shared_metrics._reset_for_tests() + relay_runtime._reset_for_tests() + + +def test_direct_runtime_records_without_enabling_a_plugin(direct_runtime, tmp_path): + base = { + "session_id": "sensitive-session", + "task_id": "task-1", + "api_request_id": "request-1", + "platform": "cli", + "provider": "custom", + "model": "gpt-sensitive-model-id", + "base_url": "http://127.0.0.1:11434/v1", + } + + assert plugins.has_hook("pre_api_request") + plugins.invoke_hook("on_session_start", **base) + plugins.invoke_hook( + "pre_api_request", + **base, + request={"body": {"messages": ["sensitive-prompt"]}}, + ) + plugins.invoke_hook( + "api_request_error", + **base, + retryable=True, + error={"message": "sensitive-error"}, + ) + plugins.invoke_hook( + "pre_api_request", + **{ + **base, + "provider": "anthropic", + "model": "claude-sonnet", + "base_url": "https://api.anthropic.com", + }, + request={"body": {"messages": ["sensitive-prompt"]}}, + ) + plugins.invoke_hook( + "post_api_request", + **{ + **base, + "provider": "anthropic", + "model": "claude-sonnet", + "base_url": "https://api.anthropic.com", + }, + response={"content": "sensitive-response"}, + ) + plugins.invoke_hook("on_session_finalize", session_id=base["session_id"]) + + starts = [event for event in direct_runtime.events if event[0] == "llm.call"] + ends = [event for event in direct_runtime.events if event[0] == "llm.call_end"] + session_starts = [ + event for event in direct_runtime.events if event[0] == "scope.push" + ] + assert len(session_starts) == 1 + assert len(starts) == 1 + assert len(ends) == 1 + assert starts[0][2] == {} + assert starts[0][3]["model_name"] == "gpt" + assert ends[0][2] == { + "call_role": "primary", + "locality": "remote", + "model_family": "claude", + "outcome": "success", + "provider_family": "direct", + } + serialized_events = json.dumps(direct_runtime.events) + assert "sensitive-prompt" not in serialized_events + assert "sensitive-response" not in serialized_events + assert "sensitive-error" not in serialized_events + assert "gpt-sensitive-model-id" not in serialized_events + assert plugins.get_plugin_manager().list_plugins() == [] + + root = tmp_path / "hermes-home" / "telemetry" / "shared_metrics" + packages = list((root / "outbox").glob("*.json")) + assert len(packages) == 1 + package = json.loads(packages[0].read_text(encoding="utf-8")) + assert package["metrics"][0]["name"] == "hermes.model_call.count" + assert package["metrics"][0]["dimensions"]["model_family"] == "claude" + assert package["metrics"][0]["value"] == 1 + + +def test_direct_runtime_is_disabled_by_default(tmp_path, monkeypatch): + fake = _Relay() + monkeypatch.setenv("HERMES_HOME", str(tmp_path / "hermes-home")) + monkeypatch.setattr(relay_runtime, "_load_nemo_relay", lambda: fake) + monkeypatch.setattr("hermes_cli.config.load_config_readonly", lambda: {}) + relay_shared_metrics._reset_for_tests() + relay_runtime._reset_for_tests() + monkeypatch.setattr(plugins, "_plugin_manager", PluginManager()) + + assert not plugins.has_hook("pre_api_request") + plugins.invoke_hook("on_session_start", session_id="s1", platform="cli") + plugins.invoke_hook("on_session_finalize", session_id="s1") + + assert fake.events == [] + assert not (tmp_path / "hermes-home" / "telemetry").exists() + relay_shared_metrics._reset_for_tests() + relay_runtime._reset_for_tests() + + +def test_core_runtime_is_fail_open_without_a_published_binding(monkeypatch, caplog): + relay_shared_metrics._reset_for_tests() + relay_runtime._reset_for_tests() + def missing_relay(name: str): + assert name == "nemo_relay" + raise ModuleNotFoundError(name) + + monkeypatch.setattr(relay_runtime.importlib, "import_module", missing_relay) + + assert relay_runtime.get_runtime() is None + assert not relay_runtime.emit_mark("hermes.probe", session_id="s1") + assert "Hermes Relay runtime initialization failed" in caplog.text + relay_runtime._reset_for_tests() + + +def test_core_mark_uses_the_shared_session_handle_without_a_plugin(direct_runtime): + plugins.invoke_hook("on_session_start", session_id="s1", platform="cli") + + handle = relay_runtime.get_session_handle("s1") + assert handle is not None + assert relay_runtime.emit_mark( + "hermes.skill.created", + session_id="s1", + data={"provenance": "agent_created"}, + metadata={"data_schema": "hermes.skill.lifecycle.v1"}, + ) + + [mark] = [event for event in direct_runtime.events if event[0] == "scope.event"] + assert mark[1] == "hermes.skill.created" + assert mark[2]["handle"] == handle + assert plugins.get_plugin_manager().list_plugins() == [] + + +def test_core_mark_lazily_starts_relay_without_metrics_or_a_plugin( + tmp_path, + monkeypatch, +): + fake = _Relay() + monkeypatch.setenv("HERMES_HOME", str(tmp_path / "hermes-home")) + monkeypatch.setattr(relay_runtime, "_load_nemo_relay", lambda: fake) + monkeypatch.setattr("hermes_cli.config.load_config_readonly", lambda: {}) + relay_shared_metrics._reset_for_tests() + relay_runtime._reset_for_tests() + monkeypatch.setattr(plugins, "_plugin_manager", PluginManager()) + + assert relay_runtime.emit_mark( + "hermes.skill.created", + session_id="s1", + data={"provenance": "agent_created"}, + ) + plugins.invoke_hook("on_session_finalize", session_id="s1") + + assert [event[0] for event in fake.events] == [ + "scope.push", + "scope.sync", + "scope.event", + "scope.sync", + "scope.pop", + "subscribers.flush", + ] + assert not any(event[0] == "subscribers.register" for event in fake.events) + assert not (tmp_path / "hermes-home" / "telemetry").exists() + relay_runtime._reset_for_tests() + + +def test_core_runtime_creates_one_session_under_concurrent_access(direct_runtime): + runtime = relay_runtime.get_runtime() + assert runtime is not None + ready = threading.Barrier(8) + sessions: list[Any] = [] + + def ensure() -> None: + ready.wait(timeout=5) + sessions.append(runtime.ensure_session({"session_id": "shared"})) + + threads = [threading.Thread(target=ensure) for _ in range(8)] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=5) + + assert all(not thread.is_alive() for thread in threads) + assert len({id(session) for session in sessions}) == 1 + assert ( + len([event for event in direct_runtime.events if event[0] == "scope.push"]) == 1 + ) + + +def test_core_runtime_parents_subagent_session_without_exposing_ids( + direct_runtime, +): + plugins.invoke_hook("on_session_start", session_id="parent", platform="cli") + parent_handle = relay_runtime.get_session_handle("parent") + + plugins.invoke_hook( + "subagent_start", + parent_session_id="parent", + child_session_id="sensitive-child", + child_subagent_id="sensitive-subagent", + ) + plugins.invoke_hook( + "on_session_start", + session_id="sensitive-child", + platform="cli", + ) + + runtime = relay_runtime.get_runtime() + assert runtime is not None + child = runtime.get_session("sensitive-child") + assert child is not None + assert child.parent_session_id == "parent" + pushes = [event for event in direct_runtime.events if event[0] == "scope.push"] + assert len(pushes) == 2 + child_kwargs = pushes[1][3] + assert child_kwargs["handle"] == parent_handle + assert child_kwargs["metadata"] == { + relay_runtime.RUNTIME_SCHEMA_KEY: relay_runtime.RUNTIME_SCHEMA_VERSION, + "nemo_relay_scope_role": "subagent", + } + assert "sensitive-child" not in json.dumps(pushes) + assert "sensitive-subagent" not in json.dumps(pushes) + + +def test_core_runtime_ignores_self_parenting_subagent_event(direct_runtime): + runtime = relay_runtime.get_runtime() + assert runtime is not None + + runtime.register_subagent( + {"parent_session_id": "same", "child_session_id": "same"} + ) + session = runtime.ensure_session({"session_id": "same"}) + + assert session is not None + assert session.parent_session_id == "" + + +def test_terminal_model_error_is_counted_as_failed(direct_runtime): + base = { + "session_id": "s1", + "task_id": "t1", + "api_request_id": "r1", + "provider": "anthropic", + "model": "claude-sonnet", + } + + plugins.invoke_hook("pre_api_request", **base) + plugins.invoke_hook("api_request_error", **base, retryable=False) + plugins.invoke_hook("on_session_finalize", session_id="s1") + + [end] = [event for event in direct_runtime.events if event[0] == "llm.call_end"] + assert end[2]["outcome"] == "failed" + + +def test_persistence_failure_does_not_escape_the_hook( + direct_runtime, + monkeypatch, + caplog, +): + runtime = relay_shared_metrics._get_runtime() + assert runtime is not None + + def fail_record(*_args: Any, **_kwargs: Any) -> None: + raise OSError("store unavailable") + + monkeypatch.setattr(runtime.subscriber.store, "record_model_call", fail_record) + plugins.invoke_hook( + "pre_api_request", + session_id="s1", + task_id="t1", + api_request_id="r1", + provider="openai", + model="gpt-5", + ) + plugins.invoke_hook( + "post_api_request", + session_id="s1", + task_id="t1", + api_request_id="r1", + provider="openai", + model="gpt-5", + ) + + assert "Unable to persist the Hermes model-call metric" in caplog.text + + +def test_close_does_not_reopen_a_session_after_scope_start_failure( + direct_runtime, + monkeypatch, +): + runtime = relay_runtime.get_runtime() + assert runtime is not None + original_push = direct_runtime.scope.push + push_attempts = 0 + + def fail_first_push(*args: Any, **kwargs: Any) -> Any: + nonlocal push_attempts + push_attempts += 1 + if push_attempts == 1: + raise RuntimeError("simulated scope failure") + return original_push(*args, **kwargs) + + direct_runtime.scope.push = fail_first_push + with pytest.raises(RuntimeError, match="simulated scope failure"): + runtime.ensure_session({"session_id": "s1"}) + + close_started = threading.Event() + allow_close = threading.Event() + original_flush = direct_runtime.subscribers.flush + + def block_flush(): + session = runtime._sessions["s1"] + assert session.closing is True + close_started.set() + assert allow_close.wait(timeout=5) + original_flush() + + direct_runtime.subscribers.flush = block_flush + close_thread = threading.Thread( + target=runtime.close_session, + args=({"session_id": "s1"},), + ) + close_thread.start() + assert close_started.wait(timeout=5) + + ensure_thread = threading.Thread( + target=runtime.ensure_session, + args=({"session_id": "s1"},), + ) + ensure_thread.start() + allow_close.set() + close_thread.join(timeout=5) + ensure_thread.join(timeout=5) + + assert not close_thread.is_alive() + assert not ensure_thread.is_alive() + assert push_attempts == 1 + assert "s1" not in runtime._sessions diff --git a/tests/plugins/test_nemo_relay_plugin.py b/tests/plugins/test_nemo_relay_plugin.py index eabb1a510b4a..c752c2fe6e0b 100644 --- a/tests/plugins/test_nemo_relay_plugin.py +++ b/tests/plugins/test_nemo_relay_plugin.py @@ -9,7 +9,6 @@ import gc import importlib import json import sys -import threading import warnings from pathlib import Path from types import SimpleNamespace @@ -17,6 +16,8 @@ from types import SimpleNamespace import pytest import yaml +from hermes_cli import plugins as plugin_api +from hermes_cli.observability import relay_runtime, relay_shared_metrics from hermes_cli.plugins import PluginManager @@ -27,10 +28,12 @@ PLUGIN_DIR = REPO_ROOT / "plugins" / "observability" / "nemo_relay" class _FakeNemoRelay: def __init__(self): self.events = [] - self._llm_handles = {} - self._subscriber_callbacks = {} - self._handle_serial = 0 - self._scope_context = contextvars.ContextVar("fake_relay_scope", default=None) + 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") self.scope = SimpleNamespace( push=self._scope_push, @@ -65,8 +68,9 @@ class _FakeNemoRelay: self.get_scope_stack = self._get_scope_stack def _scope_push(self, name, scope_type, **kwargs): - self._scope_context.set(name) - handle = ("scope", name) + 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 @@ -78,44 +82,37 @@ class _FakeNemoRelay: def _get_scope_stack(self): current = self._scope_context.get() - self.events.append(("scope_stack.sync", current)) + self.events.append(("scope.sync", current)) return current def _llm_call(self, name, request, **kwargs): - self.events.append(("llm.call.context", self._scope_context.get())) handle = ("llm", name) - if handle in self._llm_handles: - self._handle_serial += 1 - handle = ("llm", name, self._handle_serial) - self._llm_handles[handle] = kwargs + 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.context", self._scope_context.get())) self.events.append(("llm.call_end", handle, response, kwargs)) - started = self._llm_handles.pop(handle, {}) - self._emit( - _FakeEvent( - kind="scope", - category="llm", - name=handle[1], - scope_category="end", - data=response, - category_profile={"model_name": started.get("model_name")}, - metadata={ - **(started.get("metadata") or {}), - **(kwargs.get("metadata") or {}), - "otel.status_code": "OK", - }, - ) + 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)) - result = func( - _FakeLLMRequest(request.headers, {"intercepted": True, **request.content}) - ) + result = func(_FakeLLMRequest(request.headers, {"intercepted": True, **request.content})) self.events.append(("llm.execute.end", name, result, kwargs)) return result @@ -137,9 +134,7 @@ class _FakeNemoRelay: 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 - ) + return _FakeAtifExporter(self.events, session_id, agent_name, agent_version, kwargs) async def _plugin_initialize(self, config): self.events.append(("plugin.initialize", config)) @@ -152,41 +147,16 @@ class _FakeNemoRelay: self.events.append(("plugin.activate_dynamic", config, dynamic_plugins)) return _FakePluginActivation(self.events) - def _flush_subscribers(self): - self.events.append(("subscribers.flush",)) - def _register_subscriber(self, name, callback): + self._callbacks[name] = callback self.events.append(("subscribers.register", name)) - self._subscriber_callbacks[name] = callback def _deregister_subscriber(self, name): + self._callbacks.pop(name, None) self.events.append(("subscribers.deregister", name)) - return self._subscriber_callbacks.pop(name, None) is not None - def _emit(self, event): - for callback in list(self._subscriber_callbacks.values()): - callback(event) - - -class _FakeEvent: - def __init__( - self, - *, - kind, - name, - category=None, - category_profile=None, - data=None, - metadata=None, - scope_category=None, - ): - self.kind = kind - self.name = name - self.category = category - self.category_profile = category_profile - self.data = data - self.metadata = metadata - self.scope_category = scope_category + def _flush_subscribers(self): + self.events.append(("subscribers.flush",)) class _FakePluginActivation: @@ -217,20 +187,10 @@ class _FakeAtofExporter: self.config = config def register(self, name): - self.events.append(( - "atof.register", - name, - self.config.output_directory, - self.config.filename, - )) + 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, - )) + self.events.append(("atof.deregister", name, self.config.output_directory, self.config.filename)) return True @@ -251,13 +211,16 @@ class _FakeAtifExporter: 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, - }) + 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") @@ -312,25 +275,6 @@ mode = "test" return plugins_toml -def _enable_shared_metrics(tmp_path, monkeypatch) -> Path: - hermes_home = tmp_path / "hermes-home" - hermes_home.mkdir() - (hermes_home / "config.yaml").write_text( - """ -plugins: - enabled: - - observability/nemo_relay - entries: - observability/nemo_relay: - shared_metrics: - enabled: true -""", - encoding="utf-8", - ) - monkeypatch.setenv("HERMES_HOME", str(hermes_home)) - return hermes_home - - def test_manifest_fields(): data = yaml.safe_load((PLUGIN_DIR / "plugin.yaml").read_text()) assert data["name"] == "nemo_relay" @@ -374,218 +318,13 @@ def test_nemo_relay_plugin_uses_nemo_relay_runtime(monkeypatch): assert any(event[0] == "scope.push" for event in fake_relay.events) -def test_shared_metrics_default_off_does_not_create_state(tmp_path, monkeypatch): - hermes_home = tmp_path / "hermes-home" - monkeypatch.setenv("HERMES_HOME", str(hermes_home)) - fake = _FakeNemoRelay() - plugin = _fresh_plugin(monkeypatch, fake) - - plugin.on_session_start(session_id="sensitive-session", model="sensitive-model") - - assert not any(event[0] == "subscribers.register" for event in fake.events) - assert not (hermes_home / "telemetry").exists() - - -def test_shared_metrics_counts_one_logical_model_call_across_retries( - tmp_path, monkeypatch -): - hermes_home = _enable_shared_metrics(tmp_path, monkeypatch) - fake = _FakeNemoRelay() - plugin = _fresh_plugin(monkeypatch, fake) - base = { - "session_id": "sensitive-session", - "task_id": "task-1", - "api_request_id": "request-1", - "platform": "cli", - "provider": "custom", - "model": "gpt-sensitive-model-id", - "base_url": "http://127.0.0.1:11434/v1", - } - - plugin.on_session_start(**base) - plugin.on_pre_api_request( - **base, - request={"body": {"messages": ["sensitive-prompt"]}}, - ) - plugin.on_api_request_error( - **base, - retryable=True, - error={"message": "sensitive-error"}, - ) - plugin.on_pre_api_request( - **base, - request={"body": {"messages": ["sensitive-prompt"]}}, - ) - plugin.on_post_api_request( - **base, - response={"content": "sensitive-response"}, - ) - plugin.on_session_finalize(session_id=base["session_id"]) - - model_starts = [ - event for event in fake.events if event[:2] == ("llm.call", "hermes.model_call") - ] - model_ends = [ - event - for event in fake.events - if event[0] == "llm.call_end" and event[1][1] == "hermes.model_call" - ] - assert len(model_starts) == 1 - assert len(model_ends) == 1 - assert [ - event[1] - for event in fake.events - if event[0] in {"llm.call.context", "llm.call_end.context"} - ] == ["hermes.session", "hermes.session"] - assert model_starts[0][2] == {} - assert model_starts[0][3]["model_name"] == "gpt" - assert model_ends[0][2] == { - "call_role": "primary", - "locality": "local", - "outcome": "success", - "provider_family": "custom", - } - serialized_events = json.dumps(fake.events) - assert "sensitive-prompt" not in serialized_events - assert "sensitive-response" not in serialized_events - assert "sensitive-error" not in serialized_events - assert "gpt-sensitive-model-id" not in serialized_events - runtime = plugin._get_runtime() - assert runtime is not None - assert runtime.shared_metrics is not None - assert runtime.shared_metrics.store.counter_snapshot()[0]["value"] == 1 - assert ( - len( - list( - (hermes_home / "telemetry" / "shared_metrics" / "outbox").glob("*.json") - ) - ) - == 1 - ) - - -def test_shared_metrics_closes_unfinished_model_call_as_failed(tmp_path, monkeypatch): - _enable_shared_metrics(tmp_path, monkeypatch) - fake = _FakeNemoRelay() - plugin = _fresh_plugin(monkeypatch, fake) - base = { - "session_id": "s1", - "task_id": "task-1", - "api_request_id": "request-1", - "provider": "anthropic", - "model": "claude-sonnet", - } - - plugin.on_pre_api_request(**base) - plugin.on_api_request_error(**base, retryable=False, error={"message": "private"}) - plugin.on_session_end(session_id="s1", task_id="task-1", interrupted=False) - plugin.on_session_finalize(session_id="s1") - - [model_end] = [ - event - for event in fake.events - if event[0] == "llm.call_end" and event[1][1] == "hermes.model_call" - ] - assert model_end[2]["outcome"] == "failed" - runtime = plugin._get_runtime() - assert runtime is not None - assert ( - runtime.shared_metrics.store.counter_snapshot()[0]["dimensions"]["outcome"] - == "failed" - ) - - -def test_shared_metrics_persistence_failure_is_fail_open(tmp_path, monkeypatch, caplog): - _enable_shared_metrics(tmp_path, monkeypatch) - fake = _FakeNemoRelay() - plugin = _fresh_plugin(monkeypatch, fake) - runtime = plugin._get_runtime() - assert runtime is not None - assert runtime.shared_metrics is not None - - def fail_record(*_args, **_kwargs): - raise OSError("store unavailable") - - monkeypatch.setattr(runtime.shared_metrics.store, "record_model_call", fail_record) - plugin.on_pre_api_request( - session_id="s1", - task_id="t1", - api_request_id="r1", - provider="openai", - model="gpt-5", - ) - plugin.on_post_api_request( - session_id="s1", - task_id="t1", - api_request_id="r1", - provider="openai", - model="gpt-5", - ) - - assert "Unable to persist the Hermes model-call metric" in caplog.text - - -def test_shared_metrics_close_does_not_reopen_a_failed_session_scope( - tmp_path, monkeypatch -): - _enable_shared_metrics(tmp_path, monkeypatch) - fake = _FakeNemoRelay() - original_push = fake.scope.push - push_attempts = 0 - - def fail_first_metrics_scope(*args, **kwargs): - nonlocal push_attempts - push_attempts += 1 - if push_attempts == 1: - raise RuntimeError("simulated metrics scope failure") - return original_push(*args, **kwargs) - - fake.scope.push = fail_first_metrics_scope - plugin = _fresh_plugin(monkeypatch, fake) - runtime = plugin._get_runtime() - assert runtime is not None - runtime.ensure_session({"session_id": "s1"}) - state = runtime.sessions["s1"] - assert state.metrics_handle is None - - close_started = threading.Event() - allow_close = threading.Event() - original_drain = runtime._end_pending_model_calls - - def block_drain(current_state, kwargs): - assert current_state.closing is True - close_started.set() - assert allow_close.wait(timeout=5) - original_drain(current_state, kwargs) - - monkeypatch.setattr(runtime, "_end_pending_model_calls", block_drain) - close_thread = threading.Thread( - target=runtime.close_session, - args=({"session_id": "s1"},), - ) - close_thread.start() - assert close_started.wait(timeout=5) - - runtime.ensure_session({"session_id": "s1"}) - assert push_attempts == 1 - - allow_close.set() - close_thread.join(timeout=5) - assert not close_thread.is_alive() - assert "s1" not in runtime.sessions - - def test_nemo_relay_plugin_emits_llm_tool_and_exports_atif(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) monkeypatch.setenv("HERMES_NEMO_RELAY_ATOF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY", str(tmp_path / "atof") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY", str(tmp_path / "atof")) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif")) base = { "session_id": "s1", @@ -599,26 +338,15 @@ def test_nemo_relay_plugin_emits_llm_tool_and_exports_atif(tmp_path, monkeypatch api_request_id="api-1", provider="openai", model="demo-model", - request={ - "method": "POST", - "body": {"messages": [{"role": "user", "content": "hi"}]}, - }, + request={"method": "POST", "body": {"messages": [{"role": "user", "content": "hi"}]}}, ) plugin.on_post_api_request( **base, api_request_id="api-1", response={"assistant_message": {"role": "assistant", "content": "hello"}}, ) - plugin.on_pre_tool_call( - **base, tool_name="read_file", tool_call_id="tool-1", args={"path": "x"} - ) - plugin.on_post_tool_call( - **base, - tool_name="read_file", - tool_call_id="tool-1", - result='{"ok": true}', - status="ok", - ) + plugin.on_pre_tool_call(**base, tool_name="read_file", tool_call_id="tool-1", args={"path": "x"}) + plugin.on_post_tool_call(**base, tool_name="read_file", tool_call_id="tool-1", result='{"ok": true}', status="ok") plugin.on_session_end(**base, completed=True, interrupted=False) plugin.on_session_finalize(**base, reason="shutdown") @@ -633,6 +361,79 @@ def test_nemo_relay_plugin_emits_llm_tool_and_exports_atif(tmp_path, monkeypatch assert (tmp_path / "atif" / "hermes-atif-s1.json").exists() +def test_shared_metrics_and_rich_plugin_share_one_core_session( + tmp_path, + monkeypatch, +): + 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.load_config_readonly", + lambda: {"telemetry": {"shared_metrics": {"enabled": True}}}, + ) + plugin = _fresh_plugin(monkeypatch, fake) + manager = PluginManager() + manager._hooks["on_session_start"] = [plugin.on_session_start] + manager._hooks["pre_api_request"] = [plugin.on_pre_api_request] + manager._hooks["post_api_request"] = [plugin.on_post_api_request] + manager._hooks["on_session_finalize"] = [plugin.on_session_finalize] + 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", + } + plugin_api.invoke_hook("on_session_start", **event) + plugin_api.invoke_hook( + "pre_api_request", + **event, + request={"body": {"messages": [{"role": "user", "content": "hi"}]}}, + ) + plugin_api.invoke_hook( + "post_api_request", + **event, + response={"assistant_message": {"role": "assistant", "content": "hello"}}, + ) + plugin_api.invoke_hook("on_session_finalize", 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 = fake.events.index( + ("subscribers.register", "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_call.count" + assert package["metrics"][0]["value"] == 1 + assert (tmp_path / "atif" / "hermes-atif-s1.json").exists() + + def test_nemo_relay_plugin_closes_api_span_on_error(monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) @@ -660,9 +461,7 @@ def test_nemo_relay_plugin_closes_api_span_on_error(monkeypatch): call_end = next(event for event in fake.events if event[0] == "llm.call_end") assert call_end[1] == ("llm", "openai") - assert call_end[2] == { - "error": {"type": "RateLimitError", "message": "rate limited"} - } + assert call_end[2] == {"error": {"type": "RateLimitError", "message": "rate limited"}} assert call_end[3]["data"]["reason"] == "rate_limit" assert not plugin._get_runtime().sessions["s1"].llm_spans @@ -671,12 +470,8 @@ def test_nemo_relay_plugin_emits_approval_marks(monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) - plugin.on_pre_approval_request( - session_id="s1", approval_id="approval-1", tool_name="shell" - ) - plugin.on_post_approval_response( - session_id="s1", approval_id="approval-1", approved=True - ) + plugin.on_pre_approval_request(session_id="s1", approval_id="approval-1", tool_name="shell") + plugin.on_post_approval_response(session_id="s1", approval_id="approval-1", approved=True) mark_names = [event[1] for event in fake.events if event[0] == "scope.event"] assert "hermes.approval.request" in mark_names @@ -687,17 +482,13 @@ def test_nemo_relay_plugin_emits_unmatched_fallback_marks(monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) - plugin.on_post_api_request( - session_id="s1", api_request_id="missing-api", response={"ok": True} - ) + plugin.on_post_api_request(session_id="s1", api_request_id="missing-api", response={"ok": True}) plugin.on_api_request_error( session_id="s1", api_request_id="missing-api", error={"type": "TimeoutError", "message": "timed out"}, ) - plugin.on_post_tool_call( - session_id="s1", tool_call_id="missing-tool", result={"ok": True} - ) + plugin.on_post_tool_call(session_id="s1", tool_call_id="missing-tool", result={"ok": True}) mark_names = [event[1] for event in fake.events if event[0] == "scope.event"] assert "hermes.api.response.unmatched" in mark_names @@ -733,20 +524,12 @@ def test_nemo_relay_plugin_metadata_promotes_trajectory_and_subagent_ids(monkeyp telemetry_schema_version="hermes.observer.v1", ) - turn_mark = next( - event - for event in fake.events - if event[0] == "scope.event" and event[1] == "hermes.turn.start" - ) + turn_mark = next(event for event in fake.events if event[0] == "scope.event" and event[1] == "hermes.turn.start") turn_metadata = turn_mark[2]["metadata"] assert turn_metadata["session_id"] == "parent-session" assert turn_metadata["trajectory_id"] == "parent-session" - start_mark = next( - event - for event in fake.events - if event[0] == "scope.event" and event[1] == "hermes.subagent.start" - ) + start_mark = next(event for event in fake.events if event[0] == "scope.event" and event[1] == "hermes.subagent.start") start_metadata = start_mark[2]["metadata"] assert start_metadata["parent_session_id"] == "parent-session" assert start_metadata["parent_trajectory_id"] == "parent-session" @@ -755,11 +538,7 @@ def test_nemo_relay_plugin_metadata_promotes_trajectory_and_subagent_ids(monkeyp assert start_metadata["child_subagent_id"] == "child-sa" assert start_metadata["child_role"] == "leaf" - stop_mark = next( - event - for event in fake.events - if event[0] == "scope.event" and event[1] == "hermes.subagent.stop" - ) + stop_mark = next(event for event in fake.events if event[0] == "scope.event" and event[1] == "hermes.subagent.stop") assert stop_mark[2]["metadata"]["child_status"] == "completed" @@ -778,29 +557,30 @@ def test_nemo_relay_plugin_reparents_child_session_scope_for_embedded_atif(monke ) plugin.on_session_start(session_id="child-session") - child_push = next( + session_pushes = [ event for event in fake.events - if event[0] == "scope.push" and event[1] == "hermes-session-child-session" - ) + if event[0] == "scope.push" and event[1] == relay_runtime.SESSION_SCOPE + ] + assert len(session_pushes) == 2 + child_push = session_pushes[1] child_kwargs = child_push[3] - assert child_kwargs["handle"] == ("scope", "hermes-session-parent-session") + runtime = plugin._get_runtime() + assert runtime is not None + assert child_kwargs["handle"] == runtime.sessions["parent-session"].handle assert child_kwargs["metadata"]["session_id"] == "child-session" assert child_kwargs["metadata"]["trajectory_id"] == "child-session" assert child_kwargs["metadata"]["nemo_relay_scope_role"] == "subagent" assert child_kwargs["metadata"]["subagent_id"] == "child-sa" assert child_kwargs["metadata"]["parent_session_id"] == "parent-session" + assert runtime.sessions["child-session"].parent_session_id == "parent-session" -def test_nemo_relay_plugin_skips_embedded_child_atif_file_by_default( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_skips_embedded_child_atif_file_by_default(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif")) plugin.on_session_start(session_id="parent-session") plugin.on_subagent_start( @@ -818,15 +598,11 @@ def test_nemo_relay_plugin_skips_embedded_child_atif_file_by_default( assert not (tmp_path / "atif" / "hermes-atif-child-session.json").exists() -def test_nemo_relay_plugin_can_write_embedded_child_atif_file_in_all_mode( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_can_write_embedded_child_atif_file_in_all_mode(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif")) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_SUBAGENT_EXPORT_MODE", "all") plugin.on_session_start(session_id="parent-session") @@ -879,9 +655,7 @@ output_directory = "{atif_dir}" assert atif_dir.is_dir() -def test_nemo_relay_plugin_clears_plugins_toml_on_final_session_finalize_and_reinitializes( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_clears_plugins_toml_on_final_session_finalize_and_reinitializes(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) plugins_toml = tmp_path / "plugins.toml" @@ -911,9 +685,7 @@ def test_nemo_relay_plugin_activates_and_owns_dynamic_plugins(tmp_path, monkeypa plugin = _fresh_plugin(monkeypatch, fake) _enable_dynamic_plugin(tmp_path, monkeypatch) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif")) plugin.on_session_start(session_id="s1") runtime = plugin._get_runtime() @@ -941,9 +713,7 @@ def test_nemo_relay_plugin_activates_and_owns_dynamic_plugins(tmp_path, monkeypa plugin.on_session_finalize(session_id="s2", reason="shutdown") assert sum(event[0] == "plugin.activate_dynamic" for event in fake.events) == 1 - activation = next( - event for event in fake.events if event[0] == "plugin.activate_dynamic" - ) + activation = next(event for event in fake.events if event[0] == "plugin.activate_dynamic") assert "dynamic_plugins" not in activation[1] assert activation[2] == [ { @@ -960,9 +730,7 @@ def test_nemo_relay_plugin_activates_and_owns_dynamic_plugins(tmp_path, monkeypa runtime.shutdown() event_names = [event[0] for event in fake.events] - assert event_names.index("atif.deregister") < event_names.index( - "plugin.activation.close" - ) + assert event_names.index("atif.deregister") < event_names.index("plugin.activation.close") def test_nemo_relay_rejects_gateway_dynamic_config_with_actionable_diagnostic( @@ -995,9 +763,7 @@ mode = "test" assert "Use Hermes-owned [[dynamic_plugins]]" in caplog.text -def test_nemo_relay_explicit_dynamic_paths_resolve_from_plugins_toml( - tmp_path, monkeypatch -): +def test_nemo_relay_explicit_dynamic_paths_resolve_from_plugins_toml(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) config_dir = tmp_path / "config" @@ -1019,9 +785,7 @@ environment_ref = "../environments/worker-fixture" plugin.on_session_start(session_id="s1") - activation = next( - event for event in fake.events if event[0] == "plugin.activate_dynamic" - ) + activation = next(event for event in fake.events if event[0] == "plugin.activate_dynamic") assert activation[2] == [ { "plugin_id": "worker-fixture", @@ -1095,9 +859,7 @@ def test_nemo_relay_managed_llm_uses_wire_protocol_for_interceptor_dispatch( assert "rewritten_for" not in result -def test_nemo_relay_managed_llm_returns_post_next_interceptor_result( - tmp_path, monkeypatch -): +def test_nemo_relay_managed_llm_returns_post_next_interceptor_result(tmp_path, monkeypatch): fake = _FakeNemoRelay() raw_response = SimpleNamespace( model="fixture", @@ -1160,9 +922,7 @@ def test_nemo_relay_managed_tool_returns_post_interceptor_result(tmp_path, monke } -def test_nemo_relay_plugin_activates_before_registering_managed_middleware( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_activates_before_registering_managed_middleware(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) _enable_dynamic_plugin(tmp_path, monkeypatch) @@ -1255,21 +1015,16 @@ manifest_ref = "{(tmp_path / "invalid" / "relay-plugin.toml").as_posix()}" assert "no dynamic plugins will be activated" in caplog.text -def test_nemo_relay_plugin_registers_shutdown_after_dynamic_retry( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_registers_shutdown_after_dynamic_retry(tmp_path, monkeypatch): fake = _FakeNemoRelay() activation_attempts = 0 async def _flaky_activate(config, dynamic_plugins): nonlocal activation_attempts activation_attempts += 1 - fake.events.append(( - "plugin.activate_dynamic.attempt", - activation_attempts, - config, - dynamic_plugins, - )) + fake.events.append( + ("plugin.activate_dynamic.attempt", activation_attempts, config, dynamic_plugins) + ) if activation_attempts == 1: raise RuntimeError("temporary activation failure") return _FakePluginActivation(fake.events) @@ -1313,9 +1068,7 @@ def test_nemo_relay_plugin_attempts_activation_close_after_subscriber_flush_fail event_names = [event[0] for event in fake.events] assert event_names.count("subscribers.flush.failed") == 2 flush_indices = [ - index - for index, name in enumerate(event_names) - if name == "subscribers.flush.failed" + index for index, name in enumerate(event_names) if name == "subscribers.flush.failed" ] assert max(flush_indices) < event_names.index("plugin.activation.close") assert runtime._plugin_activation is None @@ -1338,9 +1091,7 @@ def test_nemo_relay_plugin_continues_shutdown_after_atif_export_failure( plugin = _fresh_plugin(monkeypatch, fake) _enable_dynamic_plugin(tmp_path, monkeypatch) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "atif")) plugin.on_session_start(session_id="s1") runtime = plugin._get_runtime() assert runtime is not None @@ -1349,19 +1100,13 @@ def test_nemo_relay_plugin_continues_shutdown_after_atif_export_failure( runtime.shutdown() event_names = [event[0] for event in fake.events] - assert event_names.index("atif.export.failed") < event_names.index( - "atif.deregister" - ) - assert event_names.index("atif.deregister") < event_names.index( - "plugin.activation.close" - ) + assert event_names.index("atif.export.failed") < event_names.index("atif.deregister") + assert event_names.index("atif.deregister") < event_names.index("plugin.activation.close") assert runtime._plugin_activation is None assert "ATIF export failed: disk full" in caplog.text -def test_nemo_relay_plugin_keeps_plugins_toml_active_while_other_sessions_remain( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_keeps_plugins_toml_active_while_other_sessions_remain(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) plugins_toml = tmp_path / "plugins.toml" @@ -1387,9 +1132,7 @@ enabled = true assert event_names.count("plugin.clear") == 1 -def test_nemo_relay_plugin_reinitializes_plugins_toml_inside_active_event_loop( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_reinitializes_plugins_toml_inside_active_event_loop(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) plugins_toml = tmp_path / "plugins.toml" @@ -1421,12 +1164,10 @@ enabled = true assert runtime is not None assert runtime._plugin_config_initialized is True scope_push_names = [event[1] for event in fake.events if event[0] == "scope.push"] - assert "hermes-session-s2" in scope_push_names + assert relay_runtime.SESSION_SCOPE in scope_push_names -def test_nemo_relay_plugin_retries_plugins_toml_after_clear_failure( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_retries_plugins_toml_after_clear_failure(tmp_path, monkeypatch): fake = _FakeNemoRelay() initialize_calls = 0 @@ -1464,12 +1205,10 @@ enabled = true assert event_names.count("plugin.initialize.attempt") == 2 assert event_names.count("plugin.clear.failed") == 1 scope_push_names = [event[1] for event in fake.events if event[0] == "scope.push"] - assert "hermes-session-s2" in scope_push_names + assert relay_runtime.SESSION_SCOPE in scope_push_names -def test_nemo_relay_plugin_disables_direct_atif_when_plugins_toml_owns_atif( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_disables_direct_atif_when_plugins_toml_owns_atif(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) plugins_toml = tmp_path / "plugins.toml" @@ -1489,9 +1228,7 @@ output_directory = "{(tmp_path / "managed-atif").as_posix()}" ) monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml)) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "direct-atif") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "direct-atif")) plugin.on_session_start(session_id="s1") plugin.on_session_finalize(session_id="s1", reason="shutdown") @@ -1503,9 +1240,7 @@ output_directory = "{(tmp_path / "managed-atif").as_posix()}" assert not (tmp_path / "direct-atif" / "hermes-atif-s1.json").exists() -def test_nemo_relay_plugin_keeps_direct_atif_when_plugins_toml_init_fails( - tmp_path, monkeypatch -): +def test_nemo_relay_plugin_keeps_direct_atif_when_plugins_toml_init_fails(tmp_path, monkeypatch): fake = _FakeNemoRelay() async def _failing_initialize(config): @@ -1531,9 +1266,7 @@ output_directory = "{(tmp_path / "managed-atif").as_posix()}" ) monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml)) monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "direct-atif") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY", str(tmp_path / "direct-atif")) plugin.on_session_start(session_id="s1") plugin.on_session_finalize(session_id="s1", reason="shutdown") @@ -1579,9 +1312,7 @@ output_directory = "{(tmp_path / "managed-atof").as_posix()}" ) monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml)) monkeypatch.setenv("HERMES_NEMO_RELAY_ATOF_ENABLED", "1") - monkeypatch.setenv( - "HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY", str(tmp_path / "direct-atof") - ) + monkeypatch.setenv("HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY", str(tmp_path / "direct-atof")) plugin.on_session_start(session_id="s1") plugin.on_session_finalize(session_id="s1", reason="shutdown") @@ -1596,9 +1327,7 @@ output_directory = "{(tmp_path / "managed-atof").as_posix()}" assert event_names.count("atof.deregister") == 1 -def test_nemo_relay_adaptive_llm_execution_middleware_preserves_raw_response( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_llm_execution_middleware_preserves_raw_response(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) plugins_toml = tmp_path / "plugins.toml" @@ -1626,9 +1355,7 @@ mode = "observe_only" SimpleNamespace( id="tool-1", type="function", - function=SimpleNamespace( - name="terminal", arguments='{"command":"pwd"}' - ), + function=SimpleNamespace(name="terminal", arguments='{"command":"pwd"}'), ) ], reasoning_content="need a tool", @@ -1662,9 +1389,7 @@ mode = "observe_only" assert response.model == "demo-model" assert response.choices == [raw_choice] assert seen_request["intercepted"] is True - execute_start = next( - event for event in fake.events if event[0] == "llm.execute.start" - ) + execute_start = next(event for event in fake.events if event[0] == "llm.execute.start") assert execute_start[3]["data"]["mode"] == "observe_only" execute_end = next(event for event in fake.events if event[0] == "llm.execute.end") assert execute_end[2] == { @@ -1686,19 +1411,13 @@ mode = "observe_only" } -def test_nemo_relay_adaptive_llm_execution_preserves_downstream_error( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_llm_execution_preserves_downstream_error(tmp_path, monkeypatch): fake = _FakeNemoRelay() def native_like_execute(name, request, func, **kwargs): fake.events.append(("llm.execute.start", name, request.content, kwargs)) try: - return func( - _FakeLLMRequest( - request.headers, {"intercepted": True, **request.content} - ) - ) + return func(_FakeLLMRequest(request.headers, {"intercepted": True, **request.content})) except Exception as exc: raise RuntimeError(f"internal error: {type(exc).__name__}: {exc}") from None @@ -1738,15 +1457,9 @@ def test_nemo_relay_adaptive_llm_execution_preserves_downstream_error_with_relay def native_like_execute(name, request, func, **kwargs): try: - return func( - _FakeLLMRequest( - request.headers, {"intercepted": True, **request.content} - ) - ) + return func(_FakeLLMRequest(request.headers, {"intercepted": True, **request.content})) except Exception as exc: - raise RuntimeError( - f"internal error: {type(exc).__name__}: {exc} (retried 3x)" - ) from None + raise RuntimeError(f"internal error: {type(exc).__name__}: {exc} (retried 3x)") from None fake.llm.execute = native_like_execute plugin = _fresh_plugin(monkeypatch, fake) @@ -1773,9 +1486,7 @@ def test_nemo_relay_adaptive_llm_execution_preserves_downstream_error_with_relay assert caught.value.status_code == 403 -def test_nemo_relay_adaptive_llm_execution_keeps_unrelated_internal_error( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_llm_execution_keeps_unrelated_internal_error(tmp_path, monkeypatch): fake = _FakeNemoRelay() relay_error = RuntimeError("internal error: relay setup failed") @@ -1803,17 +1514,11 @@ def test_nemo_relay_adaptive_llm_execution_keeps_wrapped_relay_error_after_downs tmp_path, monkeypatch ): fake = _FakeNemoRelay() - relay_error = RuntimeError( - "internal error: RuntimeError: relay policy blocked after downstream" - ) + relay_error = RuntimeError("internal error: RuntimeError: relay policy blocked after downstream") def translated_execute(name, request, func, **kwargs): try: - return func( - _FakeLLMRequest( - request.headers, {"intercepted": True, **request.content} - ) - ) + return func(_FakeLLMRequest(request.headers, {"intercepted": True, **request.content})) except Exception: raise relay_error @@ -1836,9 +1541,7 @@ def test_nemo_relay_adaptive_llm_execution_keeps_wrapped_relay_error_after_downs assert caught.value is relay_error -def test_nemo_relay_adaptive_llm_execution_keeps_relay_translated_error( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_llm_execution_keeps_relay_translated_error(tmp_path, monkeypatch): fake = _FakeNemoRelay() class RelayPolicyError(Exception): @@ -1848,11 +1551,7 @@ def test_nemo_relay_adaptive_llm_execution_keeps_relay_translated_error( def translated_execute(name, request, func, **kwargs): try: - return func( - _FakeLLMRequest( - request.headers, {"intercepted": True, **request.content} - ) - ) + return func(_FakeLLMRequest(request.headers, {"intercepted": True, **request.content})) except Exception: raise relay_error @@ -1877,9 +1576,7 @@ def test_nemo_relay_adaptive_llm_execution_keeps_relay_translated_error( assert caught.value is relay_error -def test_nemo_relay_downstream_unwrap_matches_real_middleware_wrapper_shape( - monkeypatch, -): +def test_nemo_relay_downstream_unwrap_matches_real_middleware_wrapper_shape(monkeypatch): # Regression guard against core/plugin drift. The synthetic tests above model # the downstream-error wrapper with a local class, so they keep passing even # if core middleware renames its private ``_DownstreamExecutionError`` or drops @@ -1943,9 +1640,7 @@ def _adaptive_llm_execute_mode(tmp_path, monkeypatch, plugins_toml_text: str) -> next_call=lambda request: {"raw": request}, ) - execute_start = next( - event for event in fake.events if event[0] == "llm.execute.start" - ) + execute_start = next(event for event in fake.events if event[0] == "llm.execute.start") return execute_start[3]["data"]["mode"] @@ -1969,9 +1664,7 @@ version = 1 assert mode == "observe_only" -def test_nemo_relay_adaptive_llm_execution_middleware_accepts_legacy_top_level_mode( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_llm_execution_middleware_accepts_legacy_top_level_mode(tmp_path, monkeypatch): mode = _adaptive_llm_execute_mode( tmp_path, monkeypatch, @@ -1989,9 +1682,7 @@ mode = "route" assert mode == "route" -def test_nemo_relay_adaptive_llm_execution_middleware_prefers_tool_parallelism_mode( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_llm_execution_middleware_prefers_tool_parallelism_mode(tmp_path, monkeypatch): mode = _adaptive_llm_execute_mode( tmp_path, monkeypatch, @@ -2012,9 +1703,7 @@ mode = "schedule" assert mode == "schedule" -def test_nemo_relay_llm_execution_middleware_calls_through_without_adaptive( - monkeypatch, -): +def test_nemo_relay_llm_execution_middleware_calls_through_without_adaptive(monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) @@ -2030,9 +1719,7 @@ def test_nemo_relay_llm_execution_middleware_calls_through_without_adaptive( assert not any(event[0] == "llm.execute.start" for event in fake.events) -def test_nemo_relay_adaptive_tool_execution_middleware_preserves_raw_response( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_tool_execution_middleware_preserves_raw_response(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) plugins_toml = tmp_path / "plugins.toml" @@ -2070,16 +1757,12 @@ mode = "observe_only" assert response == {"raw": True, "args": {"command": "pwd", "intercepted": True}} assert seen_args["intercepted"] is True - execute_start = next( - event for event in fake.events if event[0] == "tool.execute.start" - ) + execute_start = next(event for event in fake.events if event[0] == "tool.execute.start") assert execute_start[3]["data"]["mode"] == "observe_only" assert execute_start[3]["data"]["tool_call_id"] == "tool-1" -def test_nemo_relay_adaptive_tool_execution_preserves_downstream_error( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_tool_execution_preserves_downstream_error(tmp_path, monkeypatch): fake = _FakeNemoRelay() def native_like_execute(name, args, func, **kwargs): @@ -2113,9 +1796,7 @@ def test_nemo_relay_adaptive_tool_execution_preserves_downstream_error( assert caught.value.status_code == 403 -def test_nemo_relay_adaptive_tool_execution_keeps_unrelated_internal_error( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_tool_execution_keeps_unrelated_internal_error(tmp_path, monkeypatch): fake = _FakeNemoRelay() relay_error = RuntimeError("internal error: relay setup failed") @@ -2142,9 +1823,7 @@ def test_nemo_relay_adaptive_tool_execution_keeps_wrapped_relay_error_after_down tmp_path, monkeypatch ): fake = _FakeNemoRelay() - relay_error = RuntimeError( - "internal error: RuntimeError: relay policy blocked after downstream" - ) + relay_error = RuntimeError("internal error: RuntimeError: relay policy blocked after downstream") def translated_execute(name, args, func, **kwargs): try: @@ -2170,9 +1849,7 @@ def test_nemo_relay_adaptive_tool_execution_keeps_wrapped_relay_error_after_down assert caught.value is relay_error -def test_nemo_relay_adaptive_tool_execution_keeps_relay_translated_error( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_tool_execution_keeps_relay_translated_error(tmp_path, monkeypatch): fake = _FakeNemoRelay() class RelayPolicyError(Exception): @@ -2206,9 +1883,7 @@ def test_nemo_relay_adaptive_tool_execution_keeps_relay_translated_error( assert caught.value is relay_error -def test_nemo_relay_tool_execution_middleware_calls_through_without_adaptive( - monkeypatch, -): +def test_nemo_relay_tool_execution_middleware_calls_through_without_adaptive(monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) @@ -2223,9 +1898,7 @@ def test_nemo_relay_tool_execution_middleware_calls_through_without_adaptive( assert not any(event[0] == "tool.execute.start" for event in fake.events) -def test_nemo_relay_adaptive_execution_skips_duplicate_observer_spans( - tmp_path, monkeypatch -): +def test_nemo_relay_adaptive_execution_skips_duplicate_observer_spans(tmp_path, monkeypatch): fake = _FakeNemoRelay() plugin = _fresh_plugin(monkeypatch, fake) plugins_toml = tmp_path / "plugins.toml" @@ -2257,12 +1930,8 @@ mode = "observe_only" request={"body": {"messages": [{"role": "user", "content": "hi"}]}}, ) plugin.on_post_api_request(**base, response={"ok": True}) - plugin.on_pre_tool_call( - **base, tool_name="terminal", tool_call_id="tool-1", args={"command": "pwd"} - ) - plugin.on_post_tool_call( - **base, tool_name="terminal", tool_call_id="tool-1", result={"ok": True} - ) + plugin.on_pre_tool_call(**base, tool_name="terminal", tool_call_id="tool-1", args={"command": "pwd"}) + plugin.on_post_tool_call(**base, tool_name="terminal", tool_call_id="tool-1", result={"ok": True}) plugin.on_llm_execution_middleware( **base, diff --git a/tests/test_project_metadata.py b/tests/test_project_metadata.py index 8c0836e9059e..a29bb19d83fd 100644 --- a/tests/test_project_metadata.py +++ b/tests/test_project_metadata.py @@ -3,6 +3,8 @@ from pathlib import Path import tomllib +from packaging.requirements import Requirement + def _load_optional_dependencies(): pyproject_path = Path(__file__).resolve().parents[1] / "pyproject.toml" @@ -11,6 +13,12 @@ def _load_optional_dependencies(): return project["optional-dependencies"] +def _load_project(): + pyproject_path = Path(__file__).resolve().parents[1] / "pyproject.toml" + with pyproject_path.open("rb") as handle: + return tomllib.load(handle)["project"] + + def _load_package_data(): pyproject_path = Path(__file__).resolve().parents[1] / "pyproject.toml" with pyproject_path.open("rb") as handle: @@ -222,14 +230,40 @@ def test_feishu_extra_includes_qrcode_for_qr_login(): assert any(dep.startswith("qrcode") for dep in feishu_extra) -def test_nemo_relay_extra_uses_supported_official_distribution_range(): - optional_dependencies = _load_optional_dependencies() +def test_nemo_relay_is_a_pinned_core_dependency(): + metadata = _load_project() - assert optional_dependencies["nemo-relay"] == ["nemo-relay>=0.5,<1.0"] - assert not any( - spec == "hermes-agent[nemo-relay]" - for spec in optional_dependencies["all"] + relay_dependencies = [ + dependency + for dependency in metadata["dependencies"] + if dependency.startswith("nemo-relay==") + ] + assert len(relay_dependencies) == 1 + requirement = Requirement(relay_dependencies[0]) + assert str(requirement.specifier) == "==0.5.0" + assert requirement.marker is not None + assert requirement.marker.evaluate( + {"sys_platform": "darwin", "platform_machine": "arm64"} ) + assert requirement.marker.evaluate( + {"sys_platform": "linux", "platform_machine": "x86_64"} + ) + assert requirement.marker.evaluate( + {"sys_platform": "linux", "platform_machine": "aarch64"} + ) + assert requirement.marker.evaluate( + {"sys_platform": "win32", "platform_machine": "AMD64"} + ) + assert requirement.marker.evaluate( + {"sys_platform": "win32", "platform_machine": "ARM64"} + ) + assert not requirement.marker.evaluate( + {"sys_platform": "android", "platform_machine": "aarch64"} + ) + assert not requirement.marker.evaluate( + {"sys_platform": "darwin", "platform_machine": "x86_64"} + ) + assert "nemo-relay" not in metadata["optional-dependencies"] def test_dashboard_plugin_manifests_and_assets_are_packaged(): @@ -244,6 +278,12 @@ def test_dashboard_plugin_manifests_and_assets_are_packaged(): assert "*/dashboard/dist/**/*" in plugin_data +def test_shared_metrics_schema_is_packaged(): + package_data = _load_package_data() + + assert "observability/schemas/*.json" in package_data["hermes_cli"] + + def test_nested_bundled_plugin_metadata_is_packaged(): """Nested opt-in plugins need manifests and READMEs in wheel installs.""" package_data = _load_package_data() diff --git a/uv.lock b/uv.lock index f0daf9798e5a..b993bf9d9261 100644 --- a/uv.lock +++ b/uv.lock @@ -1528,6 +1528,7 @@ dependencies = [ { name = "httpx", extra = ["socks"] }, { name = "jinja2" }, { name = "markdown" }, + { name = "nemo-relay", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux') or (platform_machine == 'AMD64' and sys_platform == 'win32') or (platform_machine == 'ARM64' and sys_platform == 'win32')" }, { name = "openai" }, { name = "packaging" }, { name = "pathspec" }, @@ -1662,9 +1663,6 @@ mistral = [ modal = [ { name = "modal" }, ] -nemo-relay = [ - { name = "nemo-relay" }, -] parallel-web = [ { name = "parallel-web" }, ] @@ -1804,7 +1802,7 @@ requires-dist = [ { name = "microsoft-teams-apps", marker = "extra == 'teams'", specifier = "==2.0.13.4" }, { name = "mistralai", marker = "extra == 'mistral'", specifier = "==2.4.8" }, { name = "modal", marker = "extra == 'modal'", specifier = "==1.3.4" }, - { name = "nemo-relay", marker = "extra == 'nemo-relay'", specifier = ">=0.5.0,<0.6.0" }, + { name = "nemo-relay", marker = "(platform_machine == 'arm64' and sys_platform == 'darwin') or (platform_machine == 'aarch64' and sys_platform == 'linux') or (platform_machine == 'x86_64' and sys_platform == 'linux') or (platform_machine == 'AMD64' and sys_platform == 'win32') or (platform_machine == 'ARM64' and sys_platform == 'win32')", specifier = "==0.5.0" }, { name = "numpy", marker = "extra == 'voice'", specifier = "==2.4.3" }, { name = "openai", specifier = "==2.24.0" }, { name = "packaging", specifier = "==26.0" }, @@ -1853,7 +1851,7 @@ requires-dist = [ { name = "websockets", specifier = "==15.0.1" }, { name = "youtube-transcript-api", marker = "extra == 'youtube'", specifier = "==1.2.4" }, ] -provides-extras = ["anthropic", "exa", "firecrawl", "parallel-web", "fal", "edge-tts", "modal", "daytona", "hindsight", "dev", "messaging", "cron", "slack", "matrix", "wecom", "cli", "tts-premium", "voice", "pty", "honcho", "supermemory", "mem0", "vision", "mcp", "nemo-relay", "homeassistant", "sms", "teams", "computer-use", "acp", "mistral", "bedrock", "vertex", "azure-identity", "termux", "termux-all", "dingtalk", "feishu", "google", "youtube", "web", "all"] +provides-extras = ["anthropic", "exa", "firecrawl", "parallel-web", "fal", "edge-tts", "modal", "daytona", "hindsight", "dev", "messaging", "cron", "slack", "matrix", "wecom", "cli", "tts-premium", "voice", "pty", "honcho", "supermemory", "mem0", "vision", "mcp", "homeassistant", "sms", "teams", "computer-use", "acp", "mistral", "bedrock", "vertex", "azure-identity", "termux", "termux-all", "dingtalk", "feishu", "google", "youtube", "web", "all"] [[package]] name = "hf-xet"