"""Tests for native NeMo Relay plugin configuration ownership.""" from __future__ import annotations import asyncio import contextvars import json import threading from types import SimpleNamespace from typing import Any import pytest from agent import relay_runtime HOST_CONFLICT = RuntimeError( "conflict: a static plugin configuration is already active; to combine static and dynamic plugins, " "provide the static components as the base configuration to dynamic plugin activation before calling " "plugin initialization" ) EMPTY_ACTIVATION_REPORT = { "config": {"diagnostics": []}, "config_paths": [], "dynamic_plugins": [], "resolved_config": {"components": []}, } ACTIVE_ACTIVATION_REPORT = { **EMPTY_ACTIVATION_REPORT, "resolved_config": {"components": [{"kind": "observability", "enabled": True}]}, } class _FakeRelay: def __init__( self, *, initialize_error: Exception | None = None, activation_close_error: Exception | None = None, activation_report: Any = ACTIVE_ACTIVATION_REPORT, ) -> None: self.events: list[tuple[Any, ...]] = [] self.initialize_error = initialize_error self.activation_close_error = activation_close_error self.activation_report = activation_report self.initialized_from: list[str | None] = [] self.ScopeType = SimpleNamespace(Agent="agent") self.plugin = SimpleNamespace(initialize=self._initialize_plugins) self.scope = SimpleNamespace( push=self._scope_push, pop=self._scope_pop, ) self.subscribers = SimpleNamespace(flush_async=self._flush_async) def get_scope_stack(self) -> None: return None async def _initialize_plugins( self, config: dict[str, Any], additional_plugins_toml: Any = None, ) -> Any: self.events.append(("plugin.initialize", config)) self.initialized_from.append( str(additional_plugins_toml) if additional_plugins_toml is not None else None ) if self.initialize_error is not None: raise self.initialize_error relay = self class _Activation: is_active = True def __init__(self) -> None: self.report = relay.activation_report async def close(self) -> None: relay.events.append(("plugin.activation.close",)) if relay.activation_close_error is not None: raise relay.activation_close_error return _Activation() def _scope_push(self, name: str, scope_type: Any, **kwargs: Any) -> Any: handle = ("scope", name, len(self.events)) 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)) async def _flush_async(self) -> None: self.events.append(("subscribers.flush_async",)) class _ConcurrentPublicationRelay(_FakeRelay): def __init__(self) -> None: super().__init__() self.publication_finished = threading.Event() async def _flush_async(self) -> None: self.events.append(("subscribers.flush_async",)) assert await asyncio.to_thread(self.publication_finished.wait, 5) class _BehavioralFakeRelay(_FakeRelay): """Record plugin interception together with the active session stack.""" def __init__(self) -> None: super().__init__() self._scope_stack = contextvars.ContextVar( "behavioral_fake_relay_scope_stack", default=None, ) self._plugin_source: str | None = None self.tools = SimpleNamespace(request_intercepts=self._request_intercepts) def get_scope_stack(self) -> Any: return self._scope_stack.get() async def _initialize_plugins( self, config: dict[str, Any], additional_plugins_toml: Any = None, ) -> Any: activation = await super()._initialize_plugins(config, additional_plugins_toml) self._plugin_source = "explicit" return activation def _scope_push(self, name: str, scope_type: Any, **kwargs: Any) -> Any: handle = super()._scope_push(name, scope_type, **kwargs) self._scope_stack.set(handle) return handle def _request_intercepts( self, tool_name: str, args: dict[str, Any], ) -> dict[str, Any]: scope_stack = self.get_scope_stack() self.events.append( ( "tools.request_intercepts", tool_name, args, self._plugin_source, scope_stack, ) ) return { **args, "relay_plugin_source": self._plugin_source, "relay_scope_stack": scope_stack, } @pytest.fixture(autouse=True) def _reset_runtime(): relay_runtime._reset_for_tests() yield relay_runtime._reset_for_tests() @pytest.fixture def explicit_static_config(tmp_path, monkeypatch): config = tmp_path / "plugins.toml" config.write_text( "version = 1\n\n" "[[components]]\n" 'kind = "observability"\n' "enabled = true\n\n" "[components.config]\n" "version = 4\n", encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) return config @pytest.mark.parametrize( ("report", "expected"), [ (EMPTY_ACTIVATION_REPORT, False), (ACTIVE_ACTIVATION_REPORT, True), ( { **EMPTY_ACTIVATION_REPORT, "resolved_config": { "components": [{"kind": "observability", "enabled": False}], }, }, False, ), ( { **EMPTY_ACTIVATION_REPORT, "resolved_config": {"components": [{"kind": "observability"}]}, }, True, ), ( { **EMPTY_ACTIVATION_REPORT, "dynamic_plugins": [{"id": "selected", "selected": True}], }, True, ), ( { **EMPTY_ACTIVATION_REPORT, "dynamic_plugins": [{"id": "disabled", "selected": False}], }, False, ), ({"config": {"diagnostics": []}}, True), ], ) def test_activation_report_controls_managed_execution(report, expected): activation = SimpleNamespace(report=report) assert relay_runtime._activation_requires_managed_execution(activation) is expected def test_unset_config_uses_relay_discovery(monkeypatch): monkeypatch.delenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, raising=False) relay = _FakeRelay() host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.ACTIVE ) assert relay.initialized_from == [None] host.ensure_session({"session_id": "session"}) assert relay.events[0] == ("plugin.initialize", {}) assert relay.events[1][0:2] == ("scope.push", relay_runtime.SESSION_SCOPE) finally: host.shutdown() assert relay.events[-2:] == [ ("subscribers.flush_async",), ("plugin.activation.close",), ] def test_empty_ambient_discovery_does_not_enable_managed_execution(monkeypatch): monkeypatch.delenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, raising=False) relay = _FakeRelay(activation_report=EMPTY_ACTIVATION_REPORT) host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.DISABLED ) assert relay.initialized_from == [None] assert relay.events == [("plugin.initialize", {})] finally: host.shutdown() assert relay.events[-2:] == [ ("subscribers.flush_async",), ("plugin.activation.close",), ] def test_first_profile_ambient_plugin_decision_applies_to_later_profile( tmp_path, monkeypatch, ): monkeypatch.delenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, raising=False) relay = _FakeRelay() host_a = relay_runtime.RelayRuntime(relay=relay, profile_key="profile-a") config = tmp_path / "plugins.toml" config.write_text("", encoding="utf-8") monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) host_b = relay_runtime.RelayRuntime(relay=relay, profile_key="profile-b") try: assert host_a.managed_execution_enabled() assert host_b.managed_execution_enabled() assert ( host_a._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.ACTIVE ) assert ( host_b._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.ACTIVE ) assert relay.events == [("plugin.initialize", {})] assert relay.initialized_from == [None] finally: host_a.shutdown() host_b.shutdown() def test_relay_initializes_explicit_plugins_before_first_session_scope( explicit_static_config, ): relay = _FakeRelay() host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert host.managed_execution_enabled() host.ensure_session({"session_id": "session"}) assert relay.events[0] == ("plugin.initialize", {}) assert relay.events[1][0:2] == ("scope.push", relay_runtime.SESSION_SCOPE) finally: host.shutdown() def test_foreign_active_plugin_configuration_is_left_unchanged( explicit_static_config, caplog, ): relay = _FakeRelay(initialize_error=HOST_CONFLICT) with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FOREIGN ) assert relay.events == [("plugin.initialize", {})] assert ("plugin.activation.close",) not in relay.events assert any(r.levelname == "WARNING" for r in caplog.records) finally: host.shutdown() assert ("subscribers.flush_async",) not in relay.events def test_dynamic_host_conflict_is_foreign_too(explicit_static_config, caplog): relay = _FakeRelay( initialize_error=RuntimeError( "conflict: plugin configuration is owned by an active dynamic plugin host" ) ) with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FOREIGN ) assert relay.events == [("plugin.initialize", {})] assert any(r.levelname == "WARNING" for r in caplog.records) finally: host.shutdown() def test_plugin_error_containing_conflict_text_is_not_misclassified( explicit_static_config, caplog, ): relay = _FakeRelay( initialize_error=RuntimeError( "registration failed: plugin configuration is owned by an active dynamic plugin host" ) ) with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert "Hermes Relay plugin initialization failed" in caplog.text assert "already active outside Hermes native ownership" not in caplog.text finally: host.shutdown() def test_real_binding_leaves_foreign_plugin_host_unchanged( explicit_static_config, ): relay = pytest.importorskip("nemo_relay") if getattr(relay, "_native", None) is None: pytest.skip("NeMo Relay native binding is unavailable on this platform") activation = relay_runtime._resolve_plugin_awaitable( relay.plugin.initialize( {}, additional_plugins_toml=explicit_static_config, ) ) try: host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FOREIGN ) assert activation.is_active finally: host.shutdown() assert activation.is_active finally: relay_runtime._resolve_plugin_awaitable(activation.close()) def test_legacy_exporter_env_warns_without_disabling_ambient_discovery( monkeypatch, caplog, ): monkeypatch.delenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, raising=False) monkeypatch.setenv("HERMES_NEMO_RELAY_ATOF_ENABLED", "1") monkeypatch.setenv("HERMES_NEMO_RELAY_ATIF_EXPORT_TIMEOUT_S", "30") relay = _FakeRelay() with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.ACTIVE ) assert relay.events == [("plugin.initialize", {})] assert relay.initialized_from == [None] assert "HERMES_NEMO_RELAY_ATOF_ENABLED" in caplog.text assert "HERMES_NEMO_RELAY_ATIF_EXPORT_TIMEOUT_S" in caplog.text assert "standard user or system plugins.toml still applies" in caplog.text finally: host.shutdown() def test_initialization_failure_is_fail_open(explicit_static_config, caplog): relay = _FakeRelay(initialize_error=RuntimeError("rejected config")) with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert any(r.levelname == "WARNING" for r in caplog.records) finally: host.shutdown() def test_later_host_shares_initialization_failure(explicit_static_config): relay = _FakeRelay(initialize_error=RuntimeError("transient failure")) failed_host = relay_runtime.RelayRuntime(relay=relay, profile_key="failed") assert not failed_host.managed_execution_enabled() assert ( failed_host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) relay.initialize_error = None later_host = relay_runtime.RelayRuntime(relay=relay, profile_key="later") try: assert not later_host.managed_execution_enabled() assert ( later_host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert relay.events == [("plugin.initialize", {})] finally: failed_host.shutdown() later_host.shutdown() retry_host = relay_runtime.RelayRuntime(relay=relay, profile_key="retry") try: assert retry_host.managed_execution_enabled() assert ( retry_host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.ACTIVE ) assert relay.events.count(("plugin.initialize", {})) == 2 finally: retry_host.shutdown() def test_missing_explicit_config_is_failed_for_all_current_hosts( tmp_path, monkeypatch, caplog, ): missing_config = tmp_path / "missing" / "plugins.toml" monkeypatch.setenv( relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(missing_config), ) relay = _FakeRelay() with caplog.at_level("WARNING"): first_host = relay_runtime.RelayRuntime(relay=relay, profile_key="first") missing_config.parent.mkdir() missing_config.write_text("", encoding="utf-8") later_host = relay_runtime.RelayRuntime(relay=relay, profile_key="later") try: for host in (first_host, later_host): assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert relay.events == [] assert any(r.levelname == "WARNING" for r in caplog.records) finally: first_host.shutdown() later_host.shutdown() def test_malformed_explicit_config_does_not_fall_back_to_discovery( tmp_path, monkeypatch, caplog, ): config = tmp_path / "plugins.toml" config.write_text("[[components]\nkind =", encoding="utf-8") monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay() with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert relay.events == [] assert any(r.levelname == "WARNING" for r in caplog.records) finally: host.shutdown() def test_present_plugins_section_is_validated_even_when_falsey( tmp_path, monkeypatch, caplog, ): config = tmp_path / "plugins.toml" config.write_text("plugins = []", encoding="utf-8") monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay(initialize_error=ValueError("'plugins' must be a table")) with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert relay.initialized_from == [str(config)] assert "'plugins' must be a table" in caplog.text finally: host.shutdown() def test_active_log_names_every_loaded_configuration_file(explicit_static_config, caplog): report = { **ACTIVE_ACTIVATION_REPORT, "config_paths": [str(explicit_static_config), "/etc/nemo-relay/plugins.toml"], } relay = _FakeRelay(activation_report=report) with caplog.at_level("INFO"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: [line] = [r.getMessage() for r in caplog.records if "active process-wide" in r.getMessage()] assert all(path in line for path in report["config_paths"]) finally: host.shutdown() def test_two_profile_hosts_initialize_once_and_clear_after_final_shutdown( explicit_static_config, caplog, ): relay = _BehavioralFakeRelay() with caplog.at_level("INFO"): host_a = relay_runtime.RelayRuntime(relay=relay, profile_key="profile-a") host_b = relay_runtime.RelayRuntime(relay=relay, profile_key="profile-b") assert relay.events == [("plugin.initialize", {})] assert host_a.managed_execution_enabled() assert host_b.managed_execution_enabled() assert ( caplog.text.count( "The Relay plugin host is active process-wide and applies to all " "profiles hosted by this Hermes process." ) == 1 ) rewritten_a = host_a.apply_tool_request_intercepts( session_id="profile-a-session", tool_name="terminal", args={"profile": "a"}, ) rewritten_b = host_b.apply_tool_request_intercepts( session_id="profile-b-session", tool_name="terminal", args={"profile": "b"}, ) assert rewritten_a["relay_plugin_source"] == "explicit" assert rewritten_b["relay_plugin_source"] == "explicit" assert rewritten_a["relay_scope_stack"] != rewritten_b["relay_scope_stack"] host_a.shutdown() assert ("plugin.activation.close",) not in relay.events host_b.shutdown() assert relay.events[-2:] == [ ("subscribers.flush_async",), ("plugin.activation.close",), ] assert relay.events.count(("plugin.initialize", {})) == 1 assert relay.events.count(("plugin.activation.close",)) == 1 pop_index = next( index for index, event in enumerate(relay.events) if event[0] == "scope.pop" ) assert pop_index < relay.events.index(("plugin.activation.close",)) def test_plugin_initialization_inside_running_event_loop(explicit_static_config): relay = _FakeRelay() async def construct_host() -> relay_runtime.RelayRuntime: return relay_runtime.RelayRuntime(relay=relay, profile_key="profile") host = asyncio.run(construct_host()) try: assert relay.events == [("plugin.initialize", {})] assert host.managed_execution_enabled() finally: host.shutdown() def test_static_plugin_cleanup_uses_async_apis_inside_running_event_loop( explicit_static_config, ): relay = _FakeRelay() async def run_lifecycle() -> None: host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") host.shutdown() asyncio.run(run_lifecycle()) assert relay.events == [ ("plugin.initialize", {}), ("subscribers.flush_async",), ("plugin.activation.close",), ] def test_dynamic_plugins_share_owned_activation_until_final_host_shutdown( tmp_path, monkeypatch, ): config = tmp_path / ".nemo-relay" / "plugins.toml" config.parent.mkdir() config.write_text( """ version = 1 [[components]] kind = "observability" enabled = true [components.config] version = 1 [[plugins.dynamic]] manifest = "plugins/native/relay-plugin.toml" [plugins.dynamic.config] mode = "strict" [[plugins.dynamic]] manifest = "plugins/worker/relay-plugin.toml" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _BehavioralFakeRelay() host_a = relay_runtime.RelayRuntime(relay=relay, profile_key="profile-a") host_b = relay_runtime.RelayRuntime(relay=relay, profile_key="profile-b") assert host_a.managed_execution_enabled() assert host_b.managed_execution_enabled() assert relay.events == [("plugin.initialize", {})] assert relay.initialized_from == [str(config)] rewritten_a = host_a.apply_tool_request_intercepts( session_id="profile-a-session", tool_name="terminal", args={"profile": "a"}, ) rewritten_b = host_b.apply_tool_request_intercepts( session_id="profile-b-session", tool_name="terminal", args={"profile": "b"}, ) assert rewritten_a["relay_plugin_source"] == "explicit" assert rewritten_b["relay_plugin_source"] == "explicit" assert rewritten_a["relay_scope_stack"] != rewritten_b["relay_scope_stack"] host_a.shutdown() assert ("plugin.activation.close",) not in relay.events host_b.shutdown() assert relay.events[-2:] == [ ("subscribers.flush_async",), ("plugin.activation.close",), ] assert relay.events.count(("plugin.activation.close",)) == 1 pop_index = next( index for index, event in enumerate(relay.events) if event[0] == "scope.pop" ) assert pop_index < relay.events.index(("plugin.activation.close",)) def test_dynamic_activation_failure_disables_plugins( tmp_path, monkeypatch, caplog, ): config = tmp_path / "plugins.toml" config.write_text( """ [[plugins.dynamic]] manifest = "relay-plugin.toml" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay(initialize_error=RuntimeError("worker rejected config")) with caplog.at_level("INFO"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert [event[0] for event in relay.events] == ["plugin.initialize"] assert any(r.levelname == "WARNING" for r in caplog.records) assert "The Relay plugin host is active process-wide" not in caplog.text finally: host.shutdown() assert ("subscribers.flush_async",) not in relay.events assert ("plugin.activation.close",) not in relay.events def test_dynamic_activation_lifecycle_inside_running_event_loop( tmp_path, monkeypatch, ): config = tmp_path / "plugins.toml" config.write_text( """ [[plugins.dynamic]] manifest = "relay-plugin.toml" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay() async def run_lifecycle() -> None: host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") assert host.managed_execution_enabled() host.shutdown() asyncio.run(run_lifecycle()) assert [event[0] for event in relay.events] == [ "plugin.initialize", "subscribers.flush_async", "plugin.activation.close", ] def test_shutdown_defers_dynamic_unload_until_async_operation_finishes( tmp_path, monkeypatch, ): config = tmp_path / "plugins.toml" config.write_text( """ [[plugins.dynamic]] manifest = "relay-plugin.toml" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay() async def run_lifecycle() -> None: host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") session = host.ensure_session({"session_id": "session"}) assert session is not None started = asyncio.Event() finish = asyncio.Event() async def in_flight_call() -> None: relay.events.append(("operation.start",)) started.set() await finish.wait() relay.events.append(("operation.end",)) operation = asyncio.create_task( host.run_in_session_async(session, in_flight_call) ) await started.wait() host.shutdown() assert host.ensure_session({"session_id": "late-session"}) is None assert ("plugin.activation.close",) not in relay.events finish.set() await operation assert await asyncio.to_thread(host._shutdown_complete.wait, 5) asyncio.run(run_lifecycle()) assert relay.events.index(("operation.end",)) < relay.events.index( ("plugin.activation.close",) ) def test_session_close_does_not_flush_during_concurrent_managed_publication( explicit_static_config, ): relay = _ConcurrentPublicationRelay() completed = threading.Event() errors: list[BaseException] = [] async def run_lifecycle() -> None: host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") closing_session = host.ensure_session({"session_id": "closing"}) active_session = host.ensure_session({"session_id": "active"}) assert closing_session is not None assert active_session is not None publication_started = asyncio.Event() finish_publication = asyncio.Event() async def managed_publication() -> None: relay.events.append(("publication.start",)) publication_started.set() await finish_publication.wait() relay.events.append(("publication.end",)) relay.publication_finished.set() publication = asyncio.create_task( host.run_in_session_async(active_session, managed_publication) ) await publication_started.wait() host.close_session({"session_id": "closing"}) relay.events.append(("session.close.returned",)) finish_publication.set() await publication host.shutdown() assert host._shutdown_complete.is_set() def run_on_event_loop_thread() -> None: try: asyncio.run(run_lifecycle()) except BaseException as exc: errors.append(exc) finally: completed.set() event_loop_thread = threading.Thread( target=run_on_event_loop_thread, name="hermes-relay-session-close-regression", daemon=True, ) event_loop_thread.start() if not completed.wait(3): # Release a broken implementation so the test process can clean up # after reporting the same deadlock guarded in production. relay.publication_finished.set() assert completed.wait(5) pytest.fail("session close blocked the active asyncio event loop") event_loop_thread.join() assert errors == [] assert relay.events.count(("subscribers.flush_async",)) == 1 assert relay.events.index(("session.close.returned",)) < relay.events.index( ("publication.end",) ) assert relay.events.index(("publication.end",)) < relay.events.index( ("subscribers.flush_async",) ) assert relay.events[-1] == ("plugin.activation.close",) def test_failed_dynamic_teardown_retains_activation_and_blocks_replacement( tmp_path, monkeypatch, caplog, ): config = tmp_path / "plugins.toml" config.write_text( """ [[plugins.dynamic]] manifest = "relay-plugin.toml" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay(activation_close_error=RuntimeError("worker still busy")) host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") with caplog.at_level("WARNING"): host.shutdown() assert any(r.levelname == "WARNING" for r in caplog.records) activation = relay_runtime._PLUGIN_CONFIGURATION._activation assert activation is not None with caplog.at_level("WARNING"): replacement = relay_runtime.RelayRuntime( relay=relay, profile_key="replacement", ) try: assert not replacement.managed_execution_enabled() assert relay_runtime._PLUGIN_CONFIGURATION._activation is activation assert relay.events.count(("plugin.initialize", {})) == 1 assert relay.events.count(("plugin.activation.close",)) == 2 assert any(r.levelname == "WARNING" for r in caplog.records) finally: replacement.shutdown() # Relay treats a close failure as terminal; only reset the permissive # fake so this process-global fixture cannot leak into later tests. relay.activation_close_error = None relay_runtime._PLUGIN_CONFIGURATION.reset_for_tests() def test_dynamic_records_are_handed_to_relay_unparsed( tmp_path, monkeypatch, ): config = tmp_path / "plugins.toml" config.write_text( """ [[plugins.dynamic]] manifest = "relay-plugin.toml" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay() host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert host.managed_execution_enabled() assert relay.events == [("plugin.initialize", {})] assert relay.initialized_from == [str(config)] finally: host.shutdown() def test_legacy_dynamic_records_are_rejected( tmp_path, monkeypatch, caplog, ): config = tmp_path / "plugins.toml" config.write_text( """ version = 1 [[dynamic_plugins]] plugin_id = "native.policy" kind = "rust_dynamic" manifest_ref = "relay-plugin.toml" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) relay = _FakeRelay() with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert relay.events == [] assert "Hermes [[dynamic_plugins]] records are unsupported" in caplog.text assert "use Relay [[plugins.dynamic]] records" in caplog.text finally: host.shutdown() def test_real_binding_hands_dynamic_records_to_relay_and_fails_open_on_rejection( tmp_path, monkeypatch, caplog, ): relay = pytest.importorskip("nemo_relay") if getattr(relay, "_native", None) is None: pytest.skip("NeMo Relay native binding is unavailable on this platform") manifest = tmp_path / "plugins" / "relay-plugin.toml" manifest.parent.mkdir() manifest.write_text( """ manifest_version = 1 [plugin] id = "fixture.native" kind = "rust_dynamic" """.strip(), encoding="utf-8", ) config = tmp_path / "plugins.toml" config.write_text( """ version = 1 [[plugins.dynamic]] manifest = "plugins/relay-plugin.toml" [plugins.dynamic.config] mode = "strict" """.strip(), encoding="utf-8", ) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config)) assert relay_runtime._configured_plugin_inputs() == config with caplog.at_level("WARNING"): host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.FAILED ) assert "Hermes Relay plugin initialization failed" in caplog.text assert "relay-plugin.toml" in caplog.text finally: host.shutdown() relay_runtime._reset_for_tests() def test_real_binding_discovers_user_and_ignores_project_config( tmp_path, monkeypatch, ): relay = pytest.importorskip("nemo_relay") if getattr(relay, "_native", None) is None: pytest.skip("NeMo Relay native binding is unavailable on this platform") project_root = tmp_path / "project" working_directory = project_root / "workspace" config_directory = project_root / ".nemo-relay" atof_dir = tmp_path / "atof" working_directory.mkdir(parents=True) config_directory.mkdir() project_config = config_directory / "plugins.toml" project_config.write_text( f""" version = 1 [[components]] kind = "observability" enabled = true [components.config] version = 3 [components.config.atof] enabled = true [[components.config.atof.sinks]] type = "file" output_directory = "{atof_dir.as_posix()}" filename = "events.jsonl" mode = "overwrite" """.strip(), encoding="utf-8", ) xdg_config_home = tmp_path / "xdg" user_config = xdg_config_home / "nemo-relay" / "plugins.toml" user_config.parent.mkdir(parents=True) user_config.write_text("", encoding="utf-8") monkeypatch.chdir(working_directory) monkeypatch.setenv("XDG_CONFIG_HOME", str(xdg_config_home)) monkeypatch.delenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, raising=False) host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert not host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.DISABLED ) report = relay_runtime._PLUGIN_CONFIGURATION._activation.report config_paths = set(report["config_paths"]) assert str(user_config) in config_paths assert str(project_config) not in config_paths finally: host.shutdown() relay_runtime._reset_for_tests() assert not (atof_dir / "events.jsonl").exists() def test_real_binding_explicit_config_replaces_user_and_ignores_project( tmp_path, monkeypatch, ): relay = pytest.importorskip("nemo_relay") if getattr(relay, "_native", None) is None: pytest.skip("NeMo Relay native binding is unavailable on this platform") project_root = tmp_path / "project" working_directory = project_root / "workspace" config_directory = project_root / ".nemo-relay" selected_directory = tmp_path / "selected-config" project_atof_dir = tmp_path / "project-atof" selected_atof_dir = tmp_path / "selected-atof" working_directory.mkdir(parents=True) config_directory.mkdir() selected_directory.mkdir() project_config = config_directory / "plugins.toml" project_config.write_text( f""" version = 1 [[components]] kind = "observability" enabled = true [components.config] version = 3 [components.config.atof] enabled = true [[components.config.atof.sinks]] type = "file" output_directory = "{project_atof_dir.as_posix()}" filename = "events.jsonl" mode = "overwrite" """.strip(), encoding="utf-8", ) selected_config = selected_directory / "plugins.toml" selected_config.write_text( f""" version = 1 [[components]] kind = "observability" enabled = true [components.config] version = 4 [components.config.atof] enabled = true [[components.config.atof.sinks]] type = "file" output_directory = "{selected_atof_dir.as_posix()}" filename = "events.jsonl" mode = "overwrite" """.strip(), encoding="utf-8", ) xdg_config_home = tmp_path / "xdg" user_config = xdg_config_home / "nemo-relay" / "plugins.toml" user_config.parent.mkdir(parents=True) user_config.write_text("", encoding="utf-8") monkeypatch.chdir(working_directory) monkeypatch.setenv("XDG_CONFIG_HOME", str(xdg_config_home)) monkeypatch.setenv( relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(selected_config), ) host = relay_runtime.RelayRuntime(relay=relay, profile_key="profile") try: assert host.managed_execution_enabled() assert ( host._plugin_configuration_state is relay_runtime._RelayPluginConfigurationState.ACTIVE ) report = relay_runtime._PLUGIN_CONFIGURATION._activation.report config_paths = set(report["config_paths"]) assert str(selected_config) in config_paths assert str(user_config) not in config_paths assert str(project_config) not in config_paths host.ensure_session({"session_id": "native-explicit-plugins"}) finally: host.shutdown() relay_runtime._reset_for_tests() assert (selected_atof_dir / "events.jsonl").is_file() assert not (project_atof_dir / "events.jsonl").exists() def test_real_binding_keeps_two_profile_trajectories_separate_in_shared_exporters( tmp_path, monkeypatch, ): relay = pytest.importorskip("nemo_relay") if getattr(relay, "_native", None) is None: pytest.skip("NeMo Relay native binding is unavailable on this platform") from agent import relay_llm, relay_tools working_directory = tmp_path / "project" / "workspace" config_directory = tmp_path / "selected-config" atof_dir = tmp_path / "atof" atif_dir = tmp_path / "atif" working_directory.mkdir(parents=True) config_directory.mkdir() config_path = config_directory / "plugins.toml" config_path.write_text( f""" version = 1 [[components]] kind = "observability" enabled = true [components.config] version = 3 [components.config.atof] enabled = true [[components.config.atof.sinks]] type = "file" output_directory = "{atof_dir.as_posix()}" filename = "events.jsonl" mode = "overwrite" [components.config.atif] enabled = true output_directory = "{atif_dir.as_posix()}" filename_template = "trajectory-{{session_id}}.json" agent_name = "Hermes Native Test" agent_version = "test" """.strip(), encoding="utf-8", ) xdg_config_home = tmp_path / "xdg" xdg_config_home.mkdir() monkeypatch.chdir(working_directory) monkeypatch.setenv("XDG_CONFIG_HOME", str(xdg_config_home)) monkeypatch.setenv(relay_runtime.RELAY_PLUGINS_CONFIG_ENV, str(config_path)) monkeypatch.setattr(relay_runtime, "_load_nemo_relay", lambda: relay) runtime_ids: dict[str, str] = {} try: for profile in ("profile-a", "profile-b"): session_id = f"native-export-{profile}" monkeypatch.setenv("HERMES_HOME", str(tmp_path / profile)) profile_key = relay_runtime.current_profile_key() lease = relay_runtime.SESSION_COORDINATOR.acquire_conversation( profile_key=profile_key, session_id=session_id, platform="cli", model="test-model", ) assert isinstance(lease.host, relay_runtime.RelayRuntime) runtime_ids[profile] = lease.host.runtime_id turn = relay_runtime.SESSION_COORDINATOR.begin_turn( lease, turn_id=f"turn-{profile}", task_id=f"task-{profile}", ) try: assert lease.host.managed_execution_enabled() relay_llm.execute( {"model": "test-model", "messages": []}, lambda _request, profile=profile: { "id": f"response-{profile}", "model": "test-model", "choices": [ { "message": { "role": "assistant", "content": "ok", }, "finish_reason": "stop", } ], }, session_id=session_id, name="test-provider", model_name="test-model", metadata={ "api_mode": "chat_completions", "api_request_id": f"request-{profile}", }, ) relay_tools.execute( "terminal", {"command": "true"}, lambda _args: {"output": "ok"}, session_id=session_id, metadata={"tool_call_id": f"tool-{profile}"}, ) finally: relay_runtime.SESSION_COORDINATOR.end_turn( turn, outcome="success", ) relay_runtime.SESSION_COORDINATOR.release_conversation(lease) relay_runtime.SESSION_COORDINATOR.finalize_conversation( profile_key=profile_key, session_id=session_id, ) finally: relay_runtime._reset_for_tests() assert (atof_dir / "events.jsonl").is_file() atof_payload = (atof_dir / "events.jsonl").read_text(encoding="utf-8") assert all(runtime_id in atof_payload for runtime_id in runtime_ids.values()) trajectories = list(atif_dir.glob("trajectory-*.json")) assert len(trajectories) == 2 observed_runtime_ids: set[str] = set() for trajectory_path in trajectories: trajectory = json.loads(trajectory_path.read_text(encoding="utf-8")) trajectory_payload = json.dumps(trajectory) matching_runtime_ids = { runtime_id for runtime_id in runtime_ids.values() if runtime_id in trajectory_payload } assert len(matching_runtime_ids) == 1 observed_runtime_ids.update(matching_runtime_ids) observed_categories = { event["category"] for event in trajectory["extra"]["observed_events"] if event["kind"] == "scope" } assert {"agent", "llm", "tool"} <= observed_categories assert observed_runtime_ids == set(runtime_ids.values())