From 5ff03cb0c411700a1d56d9e20c89cc81198067a9 Mon Sep 17 00:00:00 2001 From: ehz0ah Date: Mon, 24 Aug 2026 23:36:27 +0800 Subject: [PATCH] fix(memory): scope OpenViking user cache to connection --- plugins/memory/openviking/README.md | 6 +- plugins/memory/openviking/__init__.py | 53 ++++-- tests/openviking_plugin/test_openviking.py | 154 ++++++++++++++++-- .../memory/test_openviking_provider.py | 4 + 4 files changed, 187 insertions(+), 30 deletions(-) diff --git a/plugins/memory/openviking/README.md b/plugins/memory/openviking/README.md index 96b489db23..5ba0393a08 100644 --- a/plugins/memory/openviking/README.md +++ b/plugins/memory/openviking/README.md @@ -96,8 +96,10 @@ Hermes sends `OPENVIKING_ACCOUNT` and `OPENVIKING_USER` as identity headers. and `mode=create`. It creates peer-scoped memory files under explicit-uid `viking://user//peers/${OPENVIKING_AGENT}/memories/...` URIs, where `` is resolved client-side from `/api/v1/system/status` (server-asserted -current user, `default` fallback) — explicit-uid URIs are canonical and work -under every OpenViking auth mode and version; the `viking://~` alias only +current user). Hermes caches a confirmed user only for the active connection. +If the probe fails, Hermes uses the configured user, or `default`, for that +operation and retries the probe later. Explicit-uid URIs are canonical and +work under every OpenViking auth mode and version; the `viking://~` alias only expands for USER/ADMIN roles, not the default dev mode. Explicit remembers do not depend on session commit extraction. diff --git a/plugins/memory/openviking/__init__.py b/plugins/memory/openviking/__init__.py index 2e293b61e8..a2f0efc6f4 100644 --- a/plugins/memory/openviking/__init__.py +++ b/plugins/memory/openviking/__init__.py @@ -105,12 +105,12 @@ _PREFERENCES_SUFFIX = "memories/preferences" _ENTITIES_SUFFIX = "memories/entities" -def _resolve_user_space(client) -> str: +def _resolve_user_space(client) -> Optional[str]: """Server-asserted current user for explicit-uid URIs. - Falls back to ``"default"`` (the trusted-mode default user, identical to - the on-disk layout ``viking/default/user/default/...``) when the probe - fails or reports no user. + Return ``None`` when the probe fails or reports no user. Callers can use a + configured fallback for that operation, but must not cache an unverified + identity because a later probe can succeed. """ try: status = client.get("/api/v1/system/status") @@ -120,9 +120,11 @@ def _resolve_user_space(client) -> str: return user except Exception: logger.debug( - "OpenViking user-space probe failed; using 'default'", exc_info=True + "OpenViking user-space probe failed; using configured fallback for " + "this operation and retrying later", + exc_info=True, ) - return "default" + return None def _user_scoped_uri(user_space: str, suffix: str) -> str: @@ -2268,9 +2270,12 @@ class OpenVikingMemoryProvider(MemoryProvider): self._agent = "" self._session_id = "" self._turn_count = 0 - # Server-asserted user space for explicit-uid URIs (#91995); resolved - # lazily from /api/v1/system/status on first use. - self._user_space_cache: Optional[str] = None + # Server-asserted user space for explicit-uid URIs (#91995). Bind the + # cached value to the client object that supplied it: /reload can swap + # endpoint, credentials, and identity on this provider instance, and + # an in-flight probe from the old client must not populate the new + # connection's cache. + self._user_space_cache: Optional[tuple[Any, str]] = None self._hermes_home = "" self._run_id = uuid.uuid4().hex self._run_lock_file: Optional[Any] = None @@ -3906,13 +3911,29 @@ class OpenVikingMemoryProvider(MemoryProvider): return f"{head}{marker}{tail}" if tail else _head_only() def _user_space(self, client=None) -> str: - """Resolve (and cache) the server-asserted user space for URIs.""" - if getattr(self, "_user_space_cache", None) is None: - active = client or getattr(self, "_client", None) - self._user_space_cache = ( - _resolve_user_space(active) if active is not None else "default" - ) - return self._user_space_cache + """Resolve the user space, caching only a confirmed client identity.""" + active = client if client is not None else getattr(self, "_client", None) + cached = getattr(self, "_user_space_cache", None) + if active is not None and cached is not None and cached[0] is active: + return cached[1] + + if active is not None: + resolved = _resolve_user_space(active) + if resolved: + # The probe can overlap a config reload. Only publish it when + # this is still the provider's active client. Old in-flight + # work can use its resolved value without contaminating the + # replacement client's cache. + if active is getattr(self, "_client", None): + self._user_space_cache = (active, resolved) + return resolved + + configured = str( + getattr(active, "_user", "") + or getattr(self, "_user", "") + or "default" + ).strip() + return configured or "default" def _user_scoped_uri(self, suffix: str, client=None) -> str: return _user_scoped_uri(self._user_space(client), suffix) diff --git a/tests/openviking_plugin/test_openviking.py b/tests/openviking_plugin/test_openviking.py index 5b8d26c31c..c7cef68122 100644 --- a/tests/openviking_plugin/test_openviking.py +++ b/tests/openviking_plugin/test_openviking.py @@ -654,10 +654,13 @@ class TestOpenVikingAutoRecallPrefetch: if parsed.path == "/health": self._send_json({"status": "ok", "healthy": True, "version": "test"}) return + if parsed.path == "/api/v1/system/status": + self._send_json({"status": "ok", "result": {"user": "user"}}) + return if parsed.path == "/api/v1/content/read": query = parse_qs(parsed.query) uri = query.get("uri", [""])[0] - if uri == "viking://user/default/memories/profile.md": + if uri == "viking://user/user/memories/profile.md": self._send_json({"result": "E2E user profile."}) return records["reads"].append(uri) @@ -667,7 +670,7 @@ class TestOpenVikingAutoRecallPrefetch: query = {key: values[0] for key, values in parse_qs(parsed.query).items()} records["listings"].append(query) uri = query.get("uri") - if uri == "viking://user/default/memories/preferences": + if uri == "viking://user/user/memories/preferences": self._send_json({ "result": [ {"isDir": True, "rel_path": "owner", "abstract": "ignored"}, @@ -679,7 +682,7 @@ class TestOpenVikingAutoRecallPrefetch: ] }) return - if uri == "viking://user/default/memories/entities": + if uri == "viking://user/user/memories/entities": self._send_json({ "result": [ { @@ -758,8 +761,8 @@ class TestOpenVikingAutoRecallPrefetch: assert "E2E abstract should not be injected." not in block assert records["reads"] == ["viking://user/peers/hermes/memories/e2e-full.md"] assert [listing["uri"] for listing in records["listings"]] == [ - "viking://user/default/memories/preferences", - "viking://user/default/memories/entities", + "viking://user/user/memories/preferences", + "viking://user/user/memories/entities", ] assert all(listing["output"] == "agent" for listing in records["listings"]) assert all(listing["recursive"].lower() == "true" for listing in records["listings"]) @@ -831,7 +834,7 @@ class TestOpenVikingMemoryUriBuilder: """URI must contain /peers/{peer_id}/ between user and memories.""" p = self._make_provider(user="alice", agent="coder") uri = p._build_memory_uri("preferences") - assert uri.startswith("viking://user/default/peers/coder/memories/preferences/mem_") + assert uri.startswith("viking://user/alice/peers/coder/memories/preferences/mem_") assert uri.endswith(".md") @@ -879,6 +882,40 @@ class TestEnsureClientReloadsEnv: assert rebuilt.api_key == "sk-fresh" assert len(constructions) == 2 + def test_rebuilt_client_resolves_its_own_user_space(self, monkeypatch): + class _StubClient: + def __init__(self, endpoint, api_key="", account="", user="", agent="hermes"): + self.endpoint = endpoint + self.api_key = api_key + self.account = account + self.user = user + self.agent = agent + + def health(self): + return True + + def get(self, path): + assert path == "/api/v1/system/status" + return {"status": "ok", "result": {"user": self.user}} + + monkeypatch.setattr("plugins.memory.openviking._VikingClient", _StubClient) + monkeypatch.setenv("OPENVIKING_ENDPOINT", "http://srv:31933") + monkeypatch.setenv("OPENVIKING_API_KEY", "") + monkeypatch.setenv("OPENVIKING_USER", "alice") + + provider = OpenVikingMemoryProvider() + provider._env_refresh_enabled = True + alice_client = provider._ensure_client() + alice_uri = provider._build_memory_uri("preferences") + + monkeypatch.setenv("OPENVIKING_USER", "bob") + bob_client = provider._ensure_client() + bob_uri = provider._build_memory_uri("preferences") + + assert bob_client is not alice_client + assert alice_uri.startswith("viking://user/alice/peers/hermes/") + assert bob_uri.startswith("viking://user/bob/peers/hermes/") + def test_handle_tool_call_reconnects_after_startup_health_failure(self, monkeypatch): instances = [] @@ -1304,30 +1341,123 @@ class TestResolveUserSpace: assert _resolve_user_space(_Client()) == "alice" - def test_probe_failure_falls_back_to_default(self): + def test_probe_failure_returns_unresolved(self): from plugins.memory.openviking import _resolve_user_space class _Client: def get(self, path): raise RuntimeError("probe down") - assert _resolve_user_space(_Client()) == "default" + assert _resolve_user_space(_Client()) is None - def test_missing_user_field_falls_back_to_default(self): + def test_missing_user_field_returns_unresolved(self): from plugins.memory.openviking import _resolve_user_space class _Client: def get(self, path): return {"status": "ok", "result": {}} - assert _resolve_user_space(_Client()) == "default" + assert _resolve_user_space(_Client()) is None class TestOpenVikingMemoryUriBuilderUserSpace: - def test_cached_user_space_flows_into_uri(self): + def test_confirmed_user_space_flows_into_uri(self): + class _Client: + def get(self, path): + assert path == "/api/v1/system/status" + return {"status": "ok", "result": {"user": "alice"}} + p = OpenVikingMemoryProvider.__new__(OpenVikingMemoryProvider) p._agent = "coder" - p._user_space_cache = "alice" + p._user = "default" + p._client = _Client() + p._user_space_cache = None uri = p._build_memory_uri("preferences") assert uri.startswith("viking://user/alice/peers/coder/memories/preferences/mem_") assert uri.endswith(".md") + + def test_transient_probe_failure_is_not_cached(self): + class _Client: + def __init__(self): + self.calls = 0 + + def get(self, path): + assert path == "/api/v1/system/status" + self.calls += 1 + if self.calls == 1: + raise RuntimeError("temporary failure") + return {"status": "ok", "result": {"user": "alice"}} + + p = OpenVikingMemoryProvider.__new__(OpenVikingMemoryProvider) + p._agent = "coder" + p._user = "fallback-user" + p._client = _Client() + p._user_space_cache = None + + first = p._build_memory_uri("preferences") + second = p._build_memory_uri("preferences") + + assert first.startswith("viking://user/fallback-user/peers/coder/") + assert second.startswith("viking://user/alice/peers/coder/") + assert p._client.calls == 2 + + def test_cache_is_bound_to_the_client_that_asserted_the_user(self): + class _Client: + def __init__(self, user): + self.user = user + self.calls = 0 + + def get(self, path): + assert path == "/api/v1/system/status" + self.calls += 1 + return {"status": "ok", "result": {"user": self.user}} + + p = OpenVikingMemoryProvider.__new__(OpenVikingMemoryProvider) + p._agent = "coder" + p._user = "default" + p._user_space_cache = None + alice = _Client("alice") + bob = _Client("bob") + + p._client = alice + assert "/user/alice/" in p._build_memory_uri("preferences") + p._client = bob + assert "/user/bob/" in p._build_memory_uri("preferences") + assert (alice.calls, bob.calls) == (1, 1) + + def test_old_inflight_probe_cannot_replace_new_client_cache(self): + old_probe_started = threading.Event() + release_old_probe = threading.Event() + + class _Client: + def __init__(self, user, wait=False): + self.user = user + self.wait = wait + + def get(self, path): + assert path == "/api/v1/system/status" + if self.wait: + old_probe_started.set() + assert release_old_probe.wait(2.0) + return {"status": "ok", "result": {"user": self.user}} + + p = OpenVikingMemoryProvider.__new__(OpenVikingMemoryProvider) + p._user = "default" + p._user_space_cache = None + alice = _Client("alice", wait=True) + bob = _Client("bob") + p._client = alice + old_result = [] + + worker = threading.Thread(target=lambda: old_result.append(p._user_space(alice))) + worker.start() + assert old_probe_started.wait(2.0) + + p._client = bob + assert p._user_space() == "bob" + release_old_probe.set() + worker.join(timeout=2.0) + + assert not worker.is_alive() + assert old_result == ["alice"] + assert p._user_space() == "bob" diff --git a/tests/plugins/memory/test_openviking_provider.py b/tests/plugins/memory/test_openviking_provider.py index b99e57f976..7492c92efb 100644 --- a/tests/plugins/memory/test_openviking_provider.py +++ b/tests/plugins/memory/test_openviking_provider.py @@ -1346,6 +1346,8 @@ def _mock_session_start_reads( request_params = dict(params or {}) uri = request_params.get("uri", "") calls.append((path, request_params, kwargs.get("timeout"))) + if path == "/api/v1/system/status": + return {"status": "ok", "result": {"user": "default"}} response = responses.get((path, uri), "") if isinstance(response, Exception): raise response @@ -1437,6 +1439,8 @@ def test_prefetch_reinjects_after_in_place_compression_same_session(): def fake_get(path, params=None, **kwargs): uri = (params or {}).get("uri", "") + if path == "/api/v1/system/status": + return {"status": "ok", "result": {"user": "default"}} if uri == "viking://user/default/memories/profile.md": return {"result": next(profiles)} return {"result": []}