diff --git a/plugins/memory/openviking/__init__.py b/plugins/memory/openviking/__init__.py index 5aef3b908503..3f7eb72c330b 100644 --- a/plugins/memory/openviking/__init__.py +++ b/plugins/memory/openviking/__init__.py @@ -1800,6 +1800,11 @@ class OpenVikingMemoryProvider(MemoryProvider): self._agent = "" self._session_id = "" self._turn_count = 0 + # Set once initialize() has resolved the connection baseline. Until then + # _ensure_client() must not re-resolve from the environment — callers + # that wire up a client directly (e.g. tests) would otherwise have it + # discarded. See _ensure_client() / #21130. + self._env_refresh_enabled = False # Guards the (_session_id, _turn_count) pair. sync_turn runs on the # MemoryManager's background sync executor while on_session_end / # on_session_switch run on the caller's thread, so the snapshot+reset @@ -2175,8 +2180,81 @@ class OpenVikingMemoryProvider(MemoryProvider): global _last_active_provider _last_active_provider = self + # Baseline established — subsequent accesses may refresh from env (#21130). + self._env_refresh_enabled = True + + def _ensure_client(self) -> Optional["_VikingClient"]: + """Return the active client, rebuilding it if the resolved config changed. + + ``/reload`` only refreshes ``os.environ`` — the existing provider + instance is not re-initialized — so OPENVIKING_* values added to + ``~/.hermes/.env`` after startup never reach the live client and tools + keep running against stale auth until the user restarts hermes (#21130). + + Re-resolve the connection settings on each access (same layering as + ``initialize``) and rebuild + health-check only when a value actually + changed; otherwise reuse the cached client so the hot path stays at one + dict comparison with zero network calls. A persistently unhealthy target + is cached as-is (``None``) under its config, so a down server does not + turn every access into a fresh health probe / restart attempt — the + background waiter started below reattaches once a local server is ready. + """ + # Before initialize() runs there is no env baseline to refresh against; + # return whatever client the caller wired up (matches legacy behavior). + if not self._env_refresh_enabled: + return self._client + + settings = _resolve_connection_settings(_load_hermes_openviking_config()) + endpoint = settings["endpoint"] + api_key = settings["api_key"] + account = settings["account"] + user = settings["user"] + agent = settings["agent"] + + config_unchanged = ( + endpoint == self._endpoint + and api_key == self._api_key + and account == self._account + and user == self._user + and agent == self._agent + ) + if config_unchanged: + # Same env as last resolve — return the cached client (or cached + # None for a known-down target) without a rebuild or network call. + return self._client + + self._endpoint = endpoint + self._api_key = api_key + self._account = account + self._user = user + self._agent = agent + + try: + client = _VikingClient( + endpoint, api_key, account=account, user=user, agent=agent, + ) + except ImportError: + logger.warning("httpx not installed — OpenViking plugin disabled") + self._client = None + return None + + health_state, health_message = _classify_runtime_openviking_health(client, endpoint) + if health_state == "healthy": + self._client = client + return self._client + if health_state == "unreachable": + # Preserve initialize()'s recovery semantics: for an unreachable + # LOCAL endpoint this (re)starts the server and attaches in the + # background; a remote one is disabled with a warning. It reads + # self._endpoint (set above) and leaves self._client = None. + self._handle_runtime_openviking_unreachable() + else: # responded but unhealthy (auth/config) + logger.warning("%s OpenViking memory disabled until config changes.", health_message) + self._client = None + return self._client + def system_prompt_block(self) -> str: - if not self._client: + if not self._ensure_client(): return "" # Provide brief info about the knowledge base try: @@ -2225,7 +2303,7 @@ class OpenVikingMemoryProvider(MemoryProvider): def prefetch(self, query: str, *, session_id: str = "") -> str: """Return recall context for this query/session.""" query_text = _derive_openviking_user_text(query).strip() - if not self._client or len(query_text) < _RECALL_QUERY_MIN_CHARS: + if len(query_text) < _RECALL_QUERY_MIN_CHARS or not self._ensure_client(): return "" effective_session_id = str(session_id or self._session_id or "").strip() @@ -3054,7 +3132,7 @@ class OpenVikingMemoryProvider(MemoryProvider): messages: Optional[List[Dict[str, Any]]] = None, ) -> None: """Record the conversation turn in OpenViking's session (non-blocking).""" - if not self._client: + if not self._ensure_client(): return user_content = _derive_openviking_user_text(user_content) @@ -3150,7 +3228,7 @@ class OpenVikingMemoryProvider(MemoryProvider): OpenViking automatically extracts 6 categories of memories: profile, preferences, entities, events, cases, and patterns. """ - if not self._client: + if not self._ensure_client(): return # Snapshot sid + turn count atomically against a concurrent sync_turn @@ -3201,7 +3279,7 @@ class OpenVikingMemoryProvider(MemoryProvider): ``_turn_count``. """ new_id = str(new_session_id or "").strip() - if not new_id or not self._client: + if not new_id or not self._ensure_client(): return rewound = bool(kwargs.get("rewound")) @@ -3252,7 +3330,7 @@ class OpenVikingMemoryProvider(MemoryProvider): metadata: Optional[Dict[str, Any]] = None, ) -> None: """Mirror successful built-in memory additions to OpenViking.""" - if not self._client or action != "add" or not content: + if action != "add" or not content or not self._ensure_client(): return subdir = _MEMORY_WRITE_TARGET_SUBDIR_MAP.get(target, _DEFAULT_MEMORY_SUBDIR) @@ -3297,7 +3375,7 @@ class OpenVikingMemoryProvider(MemoryProvider): ] def handle_tool_call(self, tool_name: str, args: dict, **kwargs) -> str: - if not self._client: + if not self._ensure_client(): return tool_error("OpenViking server not connected") try: diff --git a/tests/openviking_plugin/test_openviking.py b/tests/openviking_plugin/test_openviking.py index 70f65fa62cf8..7092acb77113 100644 --- a/tests/openviking_plugin/test_openviking.py +++ b/tests/openviking_plugin/test_openviking.py @@ -1401,3 +1401,152 @@ class TestOpenVikingMemoryUriBuilder: assert f"/memories/{subdir}/mem_" in uri, ( f"subdir '{subdir}' not placed correctly in URI: {uri}" ) + + +# =================================================================== +# Issue #21130 — OPENVIKING_* not reloaded after /reload +# =================================================================== + + +class TestEnsureClientReloadsEnv: + """Verify /reload picks up new OPENVIKING_* values without a restart (#21130).""" + + def test_ensure_client_rebuilds_when_api_key_changes(self, monkeypatch): + constructions = [] + + class _StubClient: + def __init__(self, endpoint, api_key, account="", user="", agent="hermes"): + constructions.append({"endpoint": endpoint, "api_key": api_key, + "account": account, "user": user, "agent": agent}) + self.endpoint, self.api_key = endpoint, api_key + self.account, self.user, self.agent = account, user, agent + + def health(self): + return True + + monkeypatch.setattr("plugins.memory.openviking._VikingClient", _StubClient) + monkeypatch.setenv("OPENVIKING_ENDPOINT", "http://srv:31933") + monkeypatch.setenv("OPENVIKING_API_KEY", "") + + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + first = provider._ensure_client() + assert first is not None + assert first.api_key == "" + assert len(constructions) == 1 + + # Same env on second call — must reuse cached client (no rebuild). + assert provider._ensure_client() is first + assert len(constructions) == 1 + + # Simulate /reload: env now carries the new API key. + monkeypatch.setenv("OPENVIKING_API_KEY", "sk-fresh") + rebuilt = provider._ensure_client() + assert rebuilt is not None + assert rebuilt is not first + assert rebuilt.api_key == "sk-fresh" + assert len(constructions) == 2 + + def test_ensure_client_rebuilds_when_endpoint_changes(self, monkeypatch): + builds = [] + + class _StubClient: + def __init__(self, endpoint, api_key, **kw): + builds.append(endpoint) + self.endpoint = endpoint + + def health(self): + return True + + monkeypatch.setattr("plugins.memory.openviking._VikingClient", _StubClient) + monkeypatch.setenv("OPENVIKING_ENDPOINT", "http://a") + monkeypatch.setenv("OPENVIKING_API_KEY", "key") + + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + provider._ensure_client() + provider._ensure_client() # cached + monkeypatch.setenv("OPENVIKING_ENDPOINT", "http://b") + provider._ensure_client() # rebuilds + assert builds == ["http://a", "http://b"] + + def test_ensure_client_returns_none_when_health_fails(self, monkeypatch): + class _StubClient: + def __init__(self, *a, **kw): + pass + + def health(self): + return False + + monkeypatch.setattr("plugins.memory.openviking._VikingClient", _StubClient) + monkeypatch.setenv("OPENVIKING_ENDPOINT", "http://dead") + monkeypatch.setenv("OPENVIKING_API_KEY", "") + + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + assert provider._ensure_client() is None + assert provider._client is None + + def test_handle_tool_call_uses_ensure_client(self, monkeypatch): + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + + class _StubClient: + def __init__(self, *a, **kw): + pass + + def health(self): + return False + + monkeypatch.setattr("plugins.memory.openviking._VikingClient", _StubClient) + + out = provider.handle_tool_call("viking_search", {"query": "x"}) + assert "not connected" in out.lower() + + def test_prefetch_goes_through_ensure_client(self, monkeypatch): + # Pre-turn recall must refresh from env too, not just tool calls: a + # /reload before the first recall of a turn should reach the new client. + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + calls = {"n": 0} + monkeypatch.setattr(provider, "_ensure_client", lambda: (calls.__setitem__("n", calls["n"] + 1) or None)) + # Long-enough query so the length guard doesn't short-circuit first. + assert provider.prefetch("remember my project deadlines please") == "" + assert calls["n"] == 1 + + def test_on_session_switch_goes_through_ensure_client(self, monkeypatch): + # Session rotation (/new, /branch, /resume) is a live client-gated path; + # it must refresh from env rather than read a stale snapshot. + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + calls = {"n": 0} + monkeypatch.setattr(provider, "_ensure_client", lambda: (calls.__setitem__("n", calls["n"] + 1) or None)) + provider.on_session_switch("new-sid-123") + assert calls["n"] == 1 + + def test_unreachable_local_endpoint_triggers_recovery(self, monkeypatch): + # A newly-resolved LOCAL endpoint that isn't up must preserve + # initialize()'s recovery: (re)start the server + attach in the + # background — not silently disable memory (sweeper feedback on #21130). + class _StubClient: + def __init__(self, *a, **kw): + pass + + def health(self): + return False # -> classified "unreachable" + + monkeypatch.setattr("plugins.memory.openviking._VikingClient", _StubClient) + monkeypatch.setattr("plugins.memory.openviking._is_local_openviking_url", lambda _u: True) + monkeypatch.setenv("OPENVIKING_ENDPOINT", "http://127.0.0.1:31933") + monkeypatch.setenv("OPENVIKING_API_KEY", "") + + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + recovered = {"n": 0} + monkeypatch.setattr( + provider, + "_handle_runtime_openviking_unreachable", + lambda *a, **k: recovered.__setitem__("n", recovered["n"] + 1), + ) + assert provider._ensure_client() is None + assert recovered["n"] == 1