Merge pull request #59846 from bbednarski9/bbednarski/nemo-relay-upgrade

feat(nemo-relay): nemo-relay observability version upgrade to support dynamic plugin activation
This commit is contained in:
Jeffrey Quesnelle 2026-07-14 13:15:49 -04:00 committed by GitHub
commit 77d5b2d573
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 859 additions and 54 deletions

View file

@ -41,7 +41,12 @@ class _FakeNemoRelay:
call_end=self._tool_call_end,
execute=self._tool_execute,
)
self.plugin = SimpleNamespace(initialize=self._plugin_initialize, clear=self._plugin_clear)
self.plugin = SimpleNamespace(
initialize=self._plugin_initialize,
clear=self._plugin_clear,
activate_dynamic_plugins=self._plugin_activate_dynamic,
)
self.subscribers = SimpleNamespace(flush=self._flush_subscribers)
self.LLMRequest = _FakeLLMRequest
self.AtofExporterConfig = _FakeAtofExporterConfig
self.AtofExporterMode = SimpleNamespace(Append="append", Overwrite="overwrite")
@ -100,6 +105,22 @@ class _FakeNemoRelay:
async def _plugin_clear(self):
self.events.append(("plugin.clear",))
async def _plugin_activate_dynamic(self, config, dynamic_plugins):
self.events.append(("plugin.activate_dynamic", config, dynamic_plugins))
return _FakePluginActivation(self.events)
def _flush_subscribers(self):
self.events.append(("subscribers.flush",))
class _FakePluginActivation:
def __init__(self, events):
self.events = events
self.report = {"diagnostics": []}
async def close(self):
self.events.append(("plugin.activation.close",))
class _FakeLLMRequest:
def __init__(self, headers, content):
@ -143,6 +164,7 @@ class _FakeAtifExporter:
return True
def export_json(self):
self.events.append(("atif.export", self.session_id))
return json.dumps({"session_id": self.session_id, "agent_name": self.agent_name})
@ -181,6 +203,26 @@ mode = "observe_only"
monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml))
def _enable_dynamic_plugin(tmp_path, monkeypatch) -> Path:
plugins_toml = tmp_path / "plugins.toml"
plugins_toml.write_text(
f"""
version = 1
[[dynamic_plugins]]
plugin_id = "fixture"
kind = "rust_dynamic"
manifest_ref = "{(tmp_path / "fixture" / "relay-plugin.toml").as_posix()}"
[dynamic_plugins.config]
mode = "test"
""",
encoding="utf-8",
)
monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml))
return plugins_toml
def test_manifest_fields():
data = yaml.safe_load((PLUGIN_DIR / "plugin.yaml").read_text())
assert data["name"] == "nemo_relay"
@ -508,6 +550,432 @@ enabled = true
assert event_names.count("plugin.clear") == 1
def test_nemo_relay_plugin_activates_and_owns_dynamic_plugins(tmp_path, monkeypatch):
fake = _FakeNemoRelay()
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"))
plugin.on_session_start(session_id="s1")
runtime = plugin._get_runtime()
assert runtime is not None
assert runtime._plugin_activation is not None
llm_result = plugin.on_llm_execution_middleware(
session_id="s1",
provider="openai",
model="fixture",
request={"messages": []},
next_call=lambda request: {"request": request},
)
tool_result = plugin.on_tool_execution_middleware(
session_id="s1",
tool_name="fixture-tool",
args={"value": 1},
next_call=lambda args: {"args": args},
)
assert llm_result["request"]["intercepted"] is True
assert tool_result["args"]["intercepted"] is True
plugin.on_session_finalize(session_id="s1", reason="shutdown")
assert runtime._plugin_activation is not None
assert not any(event[0] == "plugin.activation.close" for event in fake.events)
plugin.on_session_start(session_id="s2")
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")
assert "dynamic_plugins" not in activation[1]
assert activation[2] == [
{
"plugin_id": "fixture",
"kind": "rust_dynamic",
"manifest_ref": str(tmp_path / "fixture" / "relay-plugin.toml"),
"config": {"mode": "test"},
}
]
event_names = [event[0] for event in fake.events]
assert "plugin.clear" not in event_names
assert event_names.index("subscribers.flush") < event_names.index("atif.export")
assert event_names.index("atif.export") < event_names.index("atif.deregister")
runtime.shutdown()
event_names = [event[0] for event in fake.events]
assert event_names.index("atif.deregister") < event_names.index("plugin.activation.close")
def test_nemo_relay_rejects_gateway_dynamic_config_with_actionable_diagnostic(
tmp_path, monkeypatch, caplog
):
fake = _FakeNemoRelay()
plugin = _fresh_plugin(monkeypatch, fake)
plugins_toml = tmp_path / "plugins.toml"
plugins_toml.write_text(
"""
version = 1
[[plugins.dynamic]]
manifest = "plugins/fixture/relay-plugin.toml"
[plugins.dynamic.config]
mode = "test"
""",
encoding="utf-8",
)
monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml))
with caplog.at_level("ERROR"):
plugin.on_session_start(session_id="s1")
assert not any(event[0] == "plugin.activate_dynamic" for event in fake.events)
initialize = next(event for event in fake.events if event[0] == "plugin.initialize")
assert initialize[1] == {"version": 1}
assert "does not expose the CLI lifecycle resolver" in caplog.text
assert "Use Hermes-owned [[dynamic_plugins]]" in caplog.text
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"
config_dir.mkdir()
plugins_toml = config_dir / "plugins.toml"
plugins_toml.write_text(
"""
version = 1
[[dynamic_plugins]]
plugin_id = "worker-fixture"
kind = "worker"
manifest_ref = "../plugins/worker/relay-plugin.toml"
environment_ref = "../environments/worker-fixture"
""",
encoding="utf-8",
)
monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml))
plugin.on_session_start(session_id="s1")
activation = next(event for event in fake.events if event[0] == "plugin.activate_dynamic")
assert activation[2] == [
{
"plugin_id": "worker-fixture",
"kind": "worker",
"manifest_ref": str(tmp_path / "plugins" / "worker" / "relay-plugin.toml"),
"environment_ref": str(tmp_path / "environments" / "worker-fixture"),
"config": {},
}
]
runtime = plugin._get_runtime()
assert runtime is not None
runtime.shutdown()
@pytest.mark.parametrize(
("provider", "api_mode", "expected_surface", "should_rewrite"),
[
("custom", "chat_completions", "openai.chat_completions", True),
("openai-codex", "codex_responses", "openai.responses", True),
("anthropic", "anthropic_messages", "anthropic.messages", True),
("custom", "anthropic_messages", "anthropic.messages", True),
("bedrock", "bedrock_converse", "bedrock", False),
],
)
def test_nemo_relay_managed_llm_uses_wire_protocol_for_interceptor_dispatch(
tmp_path,
monkeypatch,
provider,
api_mode,
expected_surface,
should_rewrite,
):
fake = _FakeNemoRelay()
supported_surfaces = {
"anthropic.messages",
"openai.chat_completions",
"openai.responses",
}
def execute(name, request, func, **kwargs):
fake.events.append(("llm.execute.start", name, request.content, kwargs))
content = dict(request.content)
if name in supported_surfaces:
content["rewritten_for"] = name
result = func(_FakeLLMRequest(request.headers, content))
fake.events.append(("llm.execute.end", name, result, kwargs))
return result
fake.llm.execute = execute
plugin = _fresh_plugin(monkeypatch, fake)
_enable_dynamic_plugin(tmp_path, monkeypatch)
result = plugin.on_llm_execution_middleware(
session_id="s1",
provider=provider,
api_mode=api_mode,
model="fixture",
request={"messages": [{"role": "user", "content": "hi"}]},
next_call=lambda request: request,
)
execute_start = next(
event for event in fake.events if event[0] == "llm.execute.start"
)
assert execute_start[1] == expected_surface
assert execute_start[3]["metadata"]["provider"] == provider
assert execute_start[3]["metadata"]["api_mode"] == api_mode
if should_rewrite:
assert result["rewritten_for"] == expected_surface
else:
assert "rewritten_for" not in result
def test_nemo_relay_managed_llm_returns_post_next_interceptor_result(tmp_path, monkeypatch):
fake = _FakeNemoRelay()
raw_response = SimpleNamespace(
model="fixture",
choices=[
SimpleNamespace(
message=SimpleNamespace(role="assistant", content="raw", tool_calls=[]),
finish_reason="stop",
)
],
usage=None,
)
def execute(name, request, func, **kwargs):
del name, kwargs
normalized = func(_FakeLLMRequest(request.headers, request.content))
return {**normalized, "post_next_interceptor": True}
fake.llm.execute = execute
plugin = _fresh_plugin(monkeypatch, fake)
_enable_dynamic_plugin(tmp_path, monkeypatch)
result = plugin.on_llm_execution_middleware(
session_id="s1",
provider="openai",
api_mode="chat_completions",
model="fixture",
request={"messages": []},
next_call=lambda request: raw_response,
)
assert result["post_next_interceptor"] is True
assert result["assistant_message"]["content"] == "raw"
assert result is not raw_response
def test_nemo_relay_managed_tool_returns_post_interceptor_result(tmp_path, monkeypatch):
fake = _FakeNemoRelay()
def execute(name, args, func, **kwargs):
fake.events.append(("tool.execute.start", name, args, kwargs))
raw = func({"intercepted": True, **args})
result = {"compressed": True, "raw": raw}
fake.events.append(("tool.execute.end", name, result, kwargs))
return result
fake.tools.execute = execute
plugin = _fresh_plugin(monkeypatch, fake)
_enable_dynamic_plugin(tmp_path, monkeypatch)
result = plugin.on_tool_execution_middleware(
session_id="s1",
tool_name="fixture-tool",
args={"value": 1},
next_call=lambda args: {"tool_output": args},
)
assert result == {
"compressed": True,
"raw": {"tool_output": {"intercepted": True, "value": 1}},
}
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)
class _Context:
def register_hook(self, name, callback):
del name, callback
def register_middleware(self, name, callback):
del callback
fake.events.append(("hermes.register_middleware", name))
plugin.register(_Context())
event_names = [event[0] for event in fake.events]
assert event_names.index("plugin.activate_dynamic") < event_names.index(
"hermes.register_middleware"
)
runtime = plugin._get_runtime()
assert runtime is not None
runtime.shutdown()
def test_nemo_relay_plugin_degrades_to_static_config_on_relay_0_5(
tmp_path, monkeypatch, caplog
):
fake = _FakeNemoRelay()
delattr(fake.plugin, "activate_dynamic_plugins")
plugin = _fresh_plugin(monkeypatch, fake)
_enable_dynamic_plugin(tmp_path, monkeypatch)
with caplog.at_level("WARNING"):
plugin.on_session_start(session_id="s1")
initialize = next(event for event in fake.events if event[0] == "plugin.initialize")
assert "dynamic_plugins" not in initialize[1]
assert not any(event[0] == "plugin.activate_dynamic" for event in fake.events)
assert "available in NeMo Relay 0.6+" in caplog.text
def test_nemo_relay_plugin_rejects_invalid_dynamic_specs(tmp_path, monkeypatch, caplog):
fake = _FakeNemoRelay()
plugin = _fresh_plugin(monkeypatch, fake)
plugins_toml = tmp_path / "plugins.toml"
plugins_toml.write_text(
"""
version = 1
dynamic_plugins = [{ kind = "rust_dynamic", manifest_ref = "missing-id" }]
""",
encoding="utf-8",
)
monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml))
with caplog.at_level("WARNING"):
plugin.on_session_start(session_id="s1")
assert not any(event[0] == "plugin.activate_dynamic" for event in fake.events)
assert "plugin_id is required" in caplog.text
def test_nemo_relay_plugin_rejects_entire_mixed_valid_invalid_dynamic_request(
tmp_path, monkeypatch, caplog
):
fake = _FakeNemoRelay()
plugin = _fresh_plugin(monkeypatch, fake)
plugins_toml = tmp_path / "plugins.toml"
plugins_toml.write_text(
f"""
version = 1
[[dynamic_plugins]]
plugin_id = "valid-fixture"
kind = "rust_dynamic"
manifest_ref = "{(tmp_path / "valid" / "relay-plugin.toml").as_posix()}"
[[dynamic_plugins]]
kind = "worker"
manifest_ref = "{(tmp_path / "invalid" / "relay-plugin.toml").as_posix()}"
""",
encoding="utf-8",
)
monkeypatch.setenv("HERMES_NEMO_RELAY_PLUGINS_TOML", str(plugins_toml))
with caplog.at_level("WARNING"):
plugin.on_session_start(session_id="s1")
assert not any(event[0] == "plugin.activate_dynamic" for event in fake.events)
initialize = next(event for event in fake.events if event[0] == "plugin.initialize")
assert "dynamic_plugins" not in initialize[1]
assert "no dynamic plugins will be activated" in caplog.text
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)
)
if activation_attempts == 1:
raise RuntimeError("temporary activation failure")
return _FakePluginActivation(fake.events)
fake.plugin.activate_dynamic_plugins = _flaky_activate
plugin = _fresh_plugin(monkeypatch, fake)
_enable_dynamic_plugin(tmp_path, monkeypatch)
plugin.on_session_start(session_id="s1")
plugin.on_session_finalize(session_id="s1")
plugin.on_session_start(session_id="s2")
runtime = plugin._get_runtime()
assert runtime is not None
assert runtime._plugin_activation is not None
assert runtime._shutdown_registered is True
assert activation_attempts == 2
runtime.shutdown()
assert any(event[0] == "plugin.activation.close" for event in fake.events)
def test_nemo_relay_plugin_attempts_activation_close_after_subscriber_flush_failure(
tmp_path, monkeypatch, caplog
):
fake = _FakeNemoRelay()
def _failing_flush():
fake.events.append(("subscribers.flush.failed",))
raise RuntimeError("flush boom")
fake.subscribers.flush = _failing_flush
plugin = _fresh_plugin(monkeypatch, fake)
_enable_dynamic_plugin(tmp_path, monkeypatch)
plugin.on_session_start(session_id="s1")
runtime = plugin._get_runtime()
assert runtime is not None
with caplog.at_level("WARNING"):
runtime.shutdown()
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"
]
assert max(flush_indices) < event_names.index("plugin.activation.close")
assert runtime._plugin_activation is None
assert "subscriber flush failed: flush boom" in caplog.text
def test_nemo_relay_plugin_continues_shutdown_after_atif_export_failure(
tmp_path, monkeypatch, caplog
):
fake = _FakeNemoRelay()
class _FailingAtifExporter(_FakeAtifExporter):
def export_json(self):
self.events.append(("atif.export.failed", self.session_id))
raise OSError("disk full")
fake.AtifExporter = lambda session_id, agent_name, agent_version, **kwargs: (
_FailingAtifExporter(fake.events, session_id, agent_name, agent_version, kwargs)
)
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"))
plugin.on_session_start(session_id="s1")
runtime = plugin._get_runtime()
assert runtime is not None
with caplog.at_level("WARNING"):
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 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):
fake = _FakeNemoRelay()
plugin = _fresh_plugin(monkeypatch, fake)
@ -764,15 +1232,16 @@ mode = "observe_only"
),
finish_reason="tool_calls",
)
raw_response = SimpleNamespace(
id="resp-1",
model="demo-model",
choices=[raw_choice],
usage=SimpleNamespace(prompt_tokens=3, completion_tokens=5, total_tokens=8),
)
def next_call(request):
seen_request.update(request)
return SimpleNamespace(
id="resp-1",
model="demo-model",
choices=[raw_choice],
usage=SimpleNamespace(prompt_tokens=3, completion_tokens=5, total_tokens=8),
)
return raw_response
response = plugin.on_llm_execution_middleware(
session_id="s1",
@ -786,6 +1255,7 @@ mode = "observe_only"
next_call=next_call,
)
assert response is raw_response
assert response.model == "demo-model"
assert response.choices == [raw_choice]
assert seen_request["intercepted"] is True