From a98252f223b7af7b5f2e9fea48c56a930339f375 Mon Sep 17 00:00:00 2001 From: Yingliang Zhang Date: Wed, 23 Sep 2026 20:12:49 +0800 Subject: [PATCH 1/5] fix(hermes): guard client lifecycle with a leaf lock MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _client is read/written from the retain writer, the prefetch worker and the turn/tool thread with no lock. _get_client() was check-then-act and the embedded constructor takes seconds: two threads hitting a cold or just-nulled client both construct, and the loser's client is orphaned with an aiohttp session nothing ever closes ("Unclosed client session" noise). The stale-daemon retry's inline null-and-rebuild widened the window to seconds and concurrent retries clobbered each other's replacement. - _client_lock, a LEAF lock: _get_client() takes it only on the construction/retire path; the fast path stays lock-free; never held across _run_sync/operation(client) or while taking _prefetch_lock / _pending_retain_ops_lock. - _get_client(*, retire=): the retry retires the EXACT client the operation ran with — identity passed as an argument, never shared broken-client state a concurrent retry could clobber — and rebuilds exactly once under the lock (identity CAS: a sibling's fresh rebuild is returned as-is instead of being orphaned). - shutdown() retires the client under the lock BEFORE closing it (_close_client_of(client), parameterised); a concurrent _get_client() rebuilds instead of racing the close. Ported from NousResearch/hermes-agent#117236 (closed in the hindsight-move close pass). Tracked in #4662. --- hindsight-integrations/hermes/__init__.py | 62 +++++-- .../hermes/tests/test_client_lifecycle.py | 159 ++++++++++++++++++ 2 files changed, 205 insertions(+), 16 deletions(-) create mode 100644 hindsight-integrations/hermes/tests/test_client_lifecycle.py diff --git a/hindsight-integrations/hermes/__init__.py b/hindsight-integrations/hermes/__init__.py index f078939c16..f497c883c3 100644 --- a/hindsight-integrations/hermes/__init__.py +++ b/hindsight-integrations/hermes/__init__.py @@ -386,6 +386,12 @@ def backup_paths(self) -> List[str]: def __init__(self): self._config = self._api_key = self._client = None + # LEAF lock guarding _client construction/retirement: two threads racing + # a cold or just-nulled client would both build one, and the loser owns + # an aiohttp session nothing ever closes. Never held across + # _run_sync/operation(client) or while taking _prefetch_lock / + # _pending_retain_ops_lock. + self._client_lock = threading.Lock() self._api_url, self._llm_base_url, self._mode = _DEFAULT_API_URL, "", "cloud" self._timeout, self._idle_timeout = _DEFAULT_TIMEOUT, _DEFAULT_IDLE_TIMEOUT self._bank_id, self._budget, self._bank_id_template = "hermes", "mid", "" @@ -735,11 +741,28 @@ def _new_cloud_client(self): ) return Hindsight(**kwargs) - def _get_client(self): - """Return the cached Hindsight client (created once, reused).""" - if self._client is None: - self._client = self._new_embedded_client() if self._mode == "local_embedded" else self._new_cloud_client() - return self._client + def _get_client(self, *, retire=None): + """Return the cached Hindsight client (created once, reused). + + *retire* is the exact client instance the caller's failed operation ran + with — identity passed as an argument, never shared broken-client state a + concurrent retry could clobber. It is dropped and the client rebuilt + exactly once, under the lock; if a sibling already rebuilt after our + failure, that fresh client is returned as-is instead of being orphaned. + """ + if retire is None and (client := self._client) is not None: + return client + with self._client_lock: + if retire is not None: + # Identity CAS: retire only the client the operation actually ran + # with. A sibling thread may have already replaced it — its fresh + # replacement must survive our retry. + if self._client is not None and self._client is not retire: + return self._client + self._client = None + if self._client is None: + self._client = self._new_embedded_client() if self._mode == "local_embedded" else self._new_cloud_client() + return self._client def _run_sync(self, coro): """Schedule *coro* on the shared loop using the configured timeout.""" @@ -749,14 +772,17 @@ def _run_hindsight_operation(self, operation): """Run an async client operation; for local_embedded, a stale-daemon connection failure recreates the client and retries once.""" try: - return self._run_sync(operation(self._get_client())) + client = self._get_client() + return self._run_sync(operation(client)) except Exception as exc: text = f"{type(exc).__name__}: {exc}".lower() if self._mode != "local_embedded" or not any(m in text for m in _RETRIABLE_CONNECTION_MARKERS): raise logger.info("Hindsight embedded daemon appears unreachable; recreating client and retrying once: %s", exc) - self._client = None - self._client = client = self._get_client() + # Retire the exact client this operation ran with (identity as an + # argument); the loser of a concurrent retry is not closed inline + # (closing against a dead daemon can hang) — only shutdown() closes. + client = self._get_client(retire=client) return self._run_sync(operation(client)) # -- retain writer thread + server-side visibility ------------------------- @@ -1564,20 +1590,20 @@ def _flush(): self._document_id, ) - def _close_client(self) -> None: + def _close_client_of(self, client) -> None: if self._mode != "local_embedded": - self._run_sync(self._client.aclose()) + self._run_sync(client.aclose()) return # HindsightEmbedded.close() closes its sync client from this thread ("attached # to a different loop" before aiohttp releases the session): aclose the inner # client on the shared loop first, then let the wrapper clean up bookkeeping. - inner_client = getattr(self._client, "_client", None) + inner_client = getattr(client, "_client", None) if inner_client is not None and hasattr(inner_client, "aclose"): _run_sync(inner_client.aclose()) with contextlib.suppress(Exception): - self._client._client = None + client._client = None with contextlib.suppress(RuntimeError): - self._client.close() + client.close() def shutdown(self) -> None: logger.debug("Hindsight shutdown: stopping writer + waiting for background threads") @@ -1594,10 +1620,14 @@ def shutdown(self) -> None: self._retain_queue.qsize(), ) self._join_prefetch(5.0) - if self._client is not None: - with contextlib.suppress(Exception): - self._close_client() + # Retire the client under the lock BEFORE closing it, so a concurrent + # _get_client() sees None and rebuilds instead of racing the close. + with self._client_lock: + client = self._client self._client = None + if client is not None: + with contextlib.suppress(Exception): + self._close_client_of(client) # The module-global loop is intentionally NOT stopped: it's shared by every # provider in the process (one per gateway chat session); stopping it would # strand siblings' aiohttp sessions ("Unclosed client session"). Daemon diff --git a/hindsight-integrations/hermes/tests/test_client_lifecycle.py b/hindsight-integrations/hermes/tests/test_client_lifecycle.py new file mode 100644 index 0000000000..5003f61b1b --- /dev/null +++ b/hindsight-integrations/hermes/tests/test_client_lifecycle.py @@ -0,0 +1,159 @@ +"""Client lifecycle lock (#11923 in the original repo) — ``_client`` is read and +written from the retain writer thread, the background prefetch worker and the +turn/tool thread. ``_get_client()`` was check-then-act and the embedded +constructor takes seconds (runtime check, daemon spawn): two threads hitting a +cold or just-nulled client both construct, one wins, and the loser's client is +orphaned with an aiohttp session nothing ever closes. The stale-daemon retry +path (``self._client = None`` then ``_get_client()``) widened the window from +microseconds to seconds, and concurrent retries clobbered each other's +replacement. NousResearch/hermes-agent#117236.""" + +import threading +import time +import types + + +def test_client_created_once_under_concurrent_first_access(provider, monkeypatch): + """8 threads against a cold cache and a slow factory: exactly one + construction, one shared object.""" + instance, _ = provider({}) + instance._client = None + built = [] + + def _slow_build(): + time.sleep(0.2) # force the check-then-act window open + client = types.SimpleNamespace(name=f"client-{len(built)}") + built.append(client) + return client + + monkeypatch.setattr(instance, "_new_cloud_client", _slow_build) + + results = [] + + def _fetch(): + results.append(instance._get_client()) + + threads = [threading.Thread(target=_fetch, daemon=True) for _ in range(8)] + for t in threads: + t.start() + for t in threads: + t.join(timeout=10.0) + + assert len(built) == 1 + assert len(results) == 8 + assert all(client is built[0] for client in results) + + +def _op_for(broken): + def op(client): + async def _attempt(): + if client is broken: + raise RuntimeError("Cannot connect to host 127.0.0.1:8888") + return "ok" + + return _attempt() + + return op + + +def test_retry_does_not_orphan_a_sibling_client(provider, monkeypatch): + """The stale-daemon retry retires the client while sibling threads sit in + ``_get_client()``; without the lock both build and the overwritten client is + an orphan nobody closes. Exactly one replacement must ever be constructed.""" + instance, _ = provider({}) + instance._mode = "local_embedded" + built = [] + + def _build(): + time.sleep(0.2) # widen the retire-and-rebuild window + client = types.SimpleNamespace() + built.append(client) + return client + + monkeypatch.setattr(instance, "_new_embedded_client", _build) + broken = _build() + instance._client = broken + + stop = threading.Event() + + def _reader(): + while not stop.is_set(): + instance._get_client() + + readers = [threading.Thread(target=_reader, daemon=True) for _ in range(4)] + for t in readers: + t.start() + try: + assert instance._run_hindsight_operation(_op_for(broken)) == "ok" + finally: + stop.set() + for t in readers: + t.join(timeout=5.0) + + assert instance._client is not broken + orphans = [c for c in built if c is not broken and c is not instance._client] + assert orphans == [] + assert sum(c is not broken for c in built) == 1 + + +def test_concurrent_retries_retire_the_broken_client_exactly_once(provider, monkeypatch): + """Concurrent stale-daemon retries each run with the same broken client; the + identity-argument retire (``_get_client(retire=client)``) must let exactly one + of them rebuild it — never a shared broken-client flag a sibling could clobber, + never two replacements orphaning each other.""" + instance, _ = provider({}) + instance._mode = "local_embedded" + built = [] + + def _build(): + time.sleep(0.2) + client = types.SimpleNamespace() + built.append(client) + return client + + monkeypatch.setattr(instance, "_new_embedded_client", _build) + broken = _build() + instance._client = broken + + results = [] + start = threading.Barrier(5) + + def _worker(): + start.wait(timeout=5.0) + results.append(instance._run_hindsight_operation(_op_for(broken))) + + threads = [threading.Thread(target=_worker, daemon=True) for _ in range(4)] + for t in threads: + t.start() + start.wait(timeout=5.0) # release all four retries at once + for t in threads: + t.join(timeout=10.0) + + assert sorted(results) == ["ok"] * 4 + assert instance._client is not broken + # The broken client is replaced exactly once, and the lone replacement survives. + assert sum(c is not broken for c in built) == 1 + assert [c for c in built if c is not broken and c is not instance._client] == [] + + +def test_shutdown_closes_retired_client_and_allows_rebuild(provider, monkeypatch): + """shutdown() retires the client under the lock BEFORE closing it, so a + concurrent ``_get_client()`` rebuilds instead of racing the close, and a later + ``_get_client()`` must never hand back the closed object.""" + instance, _ = provider({}) + closed = [] + + class _CountingClient: + async def aclose(self): + closed.append(1) + + retired = _CountingClient() + instance._client = retired + rebuilt = types.SimpleNamespace(name="rebuilt") + monkeypatch.setattr(instance, "_new_cloud_client", lambda: rebuilt) + + instance.shutdown() + + assert closed == [1] + assert instance._client is None + assert instance._get_client() is rebuilt From 3241d8e1f2c14622cb5d443e2579b7b6dec2e6ac Mon Sep 17 00:00:00 2001 From: Yingliang Zhang Date: Wed, 23 Sep 2026 22:22:15 +0800 Subject: [PATCH 2/5] chore: sync generated files for verify-generated-files (mirrors #4678) --- hindsight-integrations/hermes/__init__.py | 4 +- .../references/developer/oracle.md | 40 ++++++++++++++----- 2 files changed, 32 insertions(+), 12 deletions(-) diff --git a/hindsight-integrations/hermes/__init__.py b/hindsight-integrations/hermes/__init__.py index f497c883c3..fa2c92af09 100644 --- a/hindsight-integrations/hermes/__init__.py +++ b/hindsight-integrations/hermes/__init__.py @@ -761,7 +761,9 @@ def _get_client(self, *, retire=None): return self._client self._client = None if self._client is None: - self._client = self._new_embedded_client() if self._mode == "local_embedded" else self._new_cloud_client() + self._client = ( + self._new_embedded_client() if self._mode == "local_embedded" else self._new_cloud_client() + ) return self._client def _run_sync(self, coro): diff --git a/skills/hindsight-docs/references/developer/oracle.md b/skills/hindsight-docs/references/developer/oracle.md index 32f82f3d58..491bb71f44 100644 --- a/skills/hindsight-docs/references/developer/oracle.md +++ b/skills/hindsight-docs/references/developer/oracle.md @@ -142,17 +142,35 @@ Example: oracle+oracledb://hindsight:s3cret@db.internal:1521/ORCLPDB1 ``` -:::warning Connection support: Easy Connect only -Hindsight builds the Oracle connection from the URL as a plain -`host:port/service_name` descriptor. **Wallet-based mTLS, TLS/TCPS, and TNS -aliases or full connect descriptors are not currently supported** by the -connection layer. In practice: - -- **Oracle Autonomous Database** and other services that require a wallet / - mTLS are not supported as-is — connect to a database reachable over a direct - `host:port/service` listener. -- The driver does not negotiate TLS itself, so secure the connection at the - network layer (private networking, VPN, or a TLS-terminating proxy). +#### Full connect descriptors and TNS aliases + +The form above is Easy Connect, which can only express `host:port/service`. When you +need more — Autonomous Database over TCPS, `retry_count`, `ssl_server_dn_match`, or a +TNS alias from a `tnsnames.ora` — leave the host empty and pass the descriptor as a +`dsn` query parameter: + +``` +oracle+oracledb://USER:PASSWORD@/?dsn=DESCRIPTOR_OR_TNS_ALIAS +``` + +The descriptor has to be percent-encoded, like the password. For example, the +Autonomous Database descriptor + +``` +(description=(retry_count=20)(address=(protocol=tcps)(port=1522)(host=adb.example.com))(connect_data=(service_name=abc_low))(security=(ssl_server_dn_match=yes))) +``` + +becomes + +```bash +export HINDSIGHT_API_DATABASE_URL="oracle+oracledb://hindsight:s3cret@/?dsn=$(python3 -c 'import urllib.parse,sys;print(urllib.parse.quote(sys.stdin.read().strip()))' <<< '(description=(retry_count=20)(address=(protocol=tcps)(port=1522)(host=adb.example.com))(connect_data=(service_name=abc_low))(security=(ssl_server_dn_match=yes)))')" +``` + +:::warning Wallet-based mTLS is still not supported +Hindsight passes no wallet directory or wallet password to the driver, so +connections that require a wallet (mTLS) do not work. TCPS with +`ssl_server_dn_match` does. Otherwise secure the connection at the network layer +(private networking, VPN, or a TLS-terminating proxy). ::: ### 3. Configure Hindsight From ac0d77b70737c0d719a321be021712a2d187cd7d Mon Sep 17 00:00:00 2001 From: Yingliang Zhang Date: Wed, 23 Sep 2026 20:09:19 +0800 Subject: [PATCH 3/5] fix(hermes): pass retain_async=True to aretain_batch in tool handler The hindsight_retain tool handler called _retain_batch without forwarding the configured self._retain_async, so tool retains silently dropped the async/sync choice and aretain_batch fell back to its own server default. retain_async is a call-level arg (never an item key), so the tool handler must pass it explicitly. Forward the configured mode, matching the auto-retain path. Ported from NousResearch/hermes-agent#60648 (closed in the hindsight-move close pass). Tracked in #4662. --- hindsight-integrations/hermes/__init__.py | 6 +++++- .../tests/test_retain_async_forwarding.py | 17 +++++++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) create mode 100644 hindsight-integrations/hermes/tests/test_retain_async_forwarding.py diff --git a/hindsight-integrations/hermes/__init__.py b/hindsight-integrations/hermes/__init__.py index 4b2e634da5..05cddb40f9 100644 --- a/hindsight-integrations/hermes/__init__.py +++ b/hindsight-integrations/hermes/__init__.py @@ -1502,7 +1502,11 @@ def _tool_retain(self, args: dict) -> str: content, context=context, tags=args.get("tags"), occurred_at=args.get("occurred_at") ) logger.debug("Tool hindsight_retain: bank=%s, content_len=%d, context=%s", self._bank_id, len(content), context) - self._retain_batch(item, bank_id=self._bank_id) + # Forward the configured retain_async mode: it is a call-level arg + # (never an item key), so omitting it here drops the async/sync choice + # and aretain_batch falls back to its own server default. Matches the + # auto-retain path. + self._retain_batch(item, bank_id=self._bank_id, retain_async=self._retain_async) logger.debug("Tool hindsight_retain: success") return "Memory stored successfully." diff --git a/hindsight-integrations/hermes/tests/test_retain_async_forwarding.py b/hindsight-integrations/hermes/tests/test_retain_async_forwarding.py new file mode 100644 index 0000000000..554b3f0547 --- /dev/null +++ b/hindsight-integrations/hermes/tests/test_retain_async_forwarding.py @@ -0,0 +1,17 @@ +"""The hindsight_retain tool handler must forward the provider's configured +retain_async mode into aretain_batch as a CALL argument (retain_async is a +call-level arg, never an item key) — otherwise tool retains silently drop the +async/sync choice and diverge from the auto-retain path +(NousResearch/hermes-agent#60648).""" + +import pytest + + +@pytest.mark.parametrize("retain_async", [True, False]) +def test_retain_tool_forwards_configured_retain_async(provider, retain_async): + instance, fake = provider({"retain_async": retain_async}) + + instance.handle_tool_call("hindsight_retain", {"content": "remember this"}) + + assert fake.retains[0]["retain_async"] is retain_async + instance.shutdown() From 5a5096fc80ad4c22d2bdd87d86290f7e769cdbd6 Mon Sep 17 00:00:00 2001 From: Yingliang Zhang Date: Wed, 23 Sep 2026 20:09:07 +0800 Subject: [PATCH 4/5] fix(hermes): fence prefetch publication to the owning session generation Ported from NousResearch/hermes-agent#64745 (closed in the hindsight-move close pass). Tracked in #4662. The background prefetch worker published its recall into the session slot unconditionally; a worker outliving on_session_switch's 3s join wrote the old session's memories into the new session's slot. queue_prefetch also spawned unbounded threads with the last finisher winning the slot. Workers now capture a slot generation at spawn, queue_prefetch bumps it and skips while a prior worker runs, on_session_switch/shutdown bump it to fence late publishers, and the publish + recall are gated on the current generation. --- hindsight-integrations/hermes/__init__.py | 34 +++- .../hermes/tests/test_prefetch.py | 158 ++++++++++++++++++ 2 files changed, 191 insertions(+), 1 deletion(-) create mode 100644 hindsight-integrations/hermes/tests/test_prefetch.py diff --git a/hindsight-integrations/hermes/__init__.py b/hindsight-integrations/hermes/__init__.py index 05cddb40f9..d161191ea8 100644 --- a/hindsight-integrations/hermes/__init__.py +++ b/hindsight-integrations/hermes/__init__.py @@ -439,6 +439,10 @@ def __init__(self): self._prefetch_result, self._prefetch_count = "", 0 self._prefetch_lock = threading.Lock() self._prefetch_thread = None + # Prefetch-slot ownership (guarded by _prefetch_lock): bumped per spawned + # worker and on session switch/shutdown, so a late or superseded worker + # whose captured generation is stale can never publish into the slot. + self._prefetch_generation = 0 self._last_recall_returned, self._last_recall_count = False, 0 self._apply_recall_settings({}) @@ -1309,15 +1313,36 @@ def queue_prefetch(self, query: str, *, session_id: str = "") -> None: if self._recall_sync or self._recall_disabled(): return + # One worker at a time (mirrors _ensure_writer): rapid turns must not + # stack concurrent recalls against one embedded daemon — the LAST + # finisher would win the slot, not the newest warm. + if self._prefetch_thread is not None and self._prefetch_thread.is_alive(): + logger.debug("Prefetch: prior worker still running; skipping this warm") + return + + with self._prefetch_lock: + self._prefetch_generation += 1 + generation = self._prefetch_generation + def _run(): # Wait (bounded, off the reply path) for the just-completed turn's # retain to be recall-visible so the warmed context includes it. if self._prefetch_waits_for_retain: self._wait_for_retains_drained(self._prefetch_retain_drain_timeout) + # A superseded/shut-down worker must not even spend a recall + # against the daemon on a result nobody will use. + with self._prefetch_lock: + if generation != self._prefetch_generation or self._shutting_down.is_set(): + return text, count = self._do_recall(query) if text: + # Fenced publish: the slot can outlive the switching session's + # 3s join (drain up to 10s + recall up to 120s), so only the + # generation that owns the slot may write it — a stale worker + # must not inject the old session's memories into the new one. with self._prefetch_lock: - self._prefetch_result, self._prefetch_count = text, count + if generation == self._prefetch_generation and not self._shutting_down.is_set(): + self._prefetch_result, self._prefetch_count = text, count self._prefetch_thread = spawn_context_thread(_run, name="hindsight-prefetch") self._prefetch_thread.start() @@ -1596,6 +1621,9 @@ def _flush(): # 2. Drain the old session's in-flight prefetch and drop its result. self._join_prefetch(3.0) with self._prefetch_lock: + # Fence out any worker still running past the join: the slot now + # belongs to the new session. + self._prefetch_generation += 1 self._prefetch_result = "" # 3. Rotate to the new session. @@ -1626,6 +1654,10 @@ def shutdown(self) -> None: logger.debug("Hindsight shutdown: stopping writer + waiting for background threads") # Stop accepting retain jobs first so late sync_turn() calls are dropped. self._shutting_down.set() + # Fence any in-flight prefetch worker: a recall finishing during or + # after the joins below must never publish into a closing provider. + with self._prefetch_lock: + self._prefetch_generation += 1 # The writer finishes in-flight work then exits on the sentinel; the # bounded join keeps shutdown predictable even if the daemon is wedged. if (writer := self._writer_thread) is not None and writer.is_alive(): diff --git a/hindsight-integrations/hermes/tests/test_prefetch.py b/hindsight-integrations/hermes/tests/test_prefetch.py new file mode 100644 index 0000000000..117d7836d7 --- /dev/null +++ b/hindsight-integrations/hermes/tests/test_prefetch.py @@ -0,0 +1,158 @@ +"""Prefetch slot fencing: the background warm worker's generation fence and +session-owner gate, keeping one session's recall out of another session's slot. + +Ported from NousResearch/hermes-agent (PRs #64745 and #117244) to the standalone +plugin; the tests reuse the recording FakeClient from ``conftest`` (no daemon). +""" + +import threading + +from conftest import FakeClient, FakeRecallResponse + + +def test_stale_worker_cannot_publish_after_switch_join_timeout(provider, monkeypatch): + """A prefetch worker can legitimately outlive on_session_switch's 3s join (10s + drain + 120s recall). Its late publish lands in the NEW session's slot; the + generation fence must drop it, and the new session's first prefetch() must + inject nothing of the old one.""" + instance, _ = provider({}) + gate = threading.Event() + entered = threading.Event() + + def _gated_recall(query): + entered.set() + gate.wait(timeout=10.0) + return "- old-session memory", 1 + + monkeypatch.setattr(instance, "_do_recall", _gated_recall) + instance.queue_prefetch("old session query") + assert entered.wait(timeout=5.0), "prefetch worker never reached the recall" + + # The worker is still gated when the switch's 3.0s join times out, so the + # switch returns with the worker alive — the leaked window upstream. + instance.on_session_switch("new-sid") + gate.set() + instance._prefetch_thread.join(timeout=5.0) + + assert instance._prefetch_result == "" + assert instance.prefetch("anything") == "" + instance.shutdown() + + +def test_superseded_worker_cannot_overwrite_newer_result(provider, monkeypatch): + """Two overlapping workers: the older finishing after the newer must not + clobber the newer result sitting in the slot.""" + instance, _ = provider({}) + gate_first = threading.Event() + entered_first = threading.Event() + + def _routed_recall(query): + if query == "first query": + entered_first.set() + gate_first.wait(timeout=10.0) + return "- first (older) memory", 1 + return "- second (newer) memory", 1 + + monkeypatch.setattr(instance, "_do_recall", _routed_recall) + + instance.queue_prefetch("first query") + worker_first = instance._prefetch_thread + # Pin the overlap: the older worker must be mid-recall before the newer + # spawns, so its final publish really exercises the fence (entry barrier + # added over the upstream port to keep the interleaving deterministic). + assert entered_first.wait(timeout=5.0), "first worker never reached the recall" + # Simulate the lost-thread-handle hazard: upstream overwrote + # _prefetch_thread freely, so a dead handle lets the next warm spawn. + instance._prefetch_thread = None + instance.queue_prefetch("second query") + worker_second = instance._prefetch_thread + + # The newer worker finishes first and wins the slot. + worker_second.join(timeout=5.0) + with instance._prefetch_lock: + assert instance._prefetch_result == "- second (newer) memory" + + # The older worker finishes LAST; its result must be fenced out. + gate_first.set() + worker_first.join(timeout=5.0) + with instance._prefetch_lock: + assert instance._prefetch_result == "- second (newer) memory" + instance.shutdown() + + +def test_shutdown_fences_inflight_prefetch_publish(provider): + """A recall resuming after shutdown() began must not publish, and the fence + must prevent a recall against a client shutdown already closed.""" + instance, fake = provider({}) + gate = threading.Event() + entered = threading.Event() + + async def _gated_arecall(**kwargs): + fake.recalls.append(kwargs) + entered.set() + gate.wait(timeout=10.0) + return FakeRecallResponse(["late memory"]) + + fake.arecall = _gated_arecall + + instance.queue_prefetch("q") + assert entered.wait(timeout=5.0), "prefetch worker never reached the recall" + + # Release the recall the moment shutdown starts so shutdown's prefetch + # join returns promptly instead of burning its 5s budget. + def _release_on_shutdown(): + instance._shutting_down.wait(timeout=10.0) + gate.set() + + threading.Thread(target=_release_on_shutdown, daemon=True).start() + instance.shutdown() + instance._prefetch_thread.join(timeout=5.0) + + assert instance._prefetch_result == "" + assert instance._client is None + assert instance.prefetch("q") == "" + + +def test_queue_prefetch_skips_while_prior_worker_running(provider): + """Rapid turns must warm serially: while one worker runs, further + queue_prefetch calls neither spawn nor bump — only ONE recall happens.""" + instance, fake = provider({}) + gate = threading.Event() + entered = threading.Event() + + async def _gated_arecall(**kwargs): + fake.recalls.append(kwargs) + entered.set() + gate.wait(timeout=10.0) + return FakeRecallResponse(["m"]) + + fake.arecall = _gated_arecall + + instance.queue_prefetch("q1") + assert entered.wait(timeout=5.0), "prefetch worker never reached the recall" + first = instance._prefetch_thread + + instance.queue_prefetch("q2") + instance.queue_prefetch("q3") + + gate.set() + first.join(timeout=5.0) + + # The live worker was never replaced and nothing else ever recalled. + assert instance._prefetch_thread is first + assert len(fake.recalls) == 1 + instance.shutdown() + + +def test_same_session_boundary_still_publishes(provider): + """Positive control: with no session switch the fence must be a no-op — an + uninterrupted warm → prefetch cycle delivers exactly as upstream.""" + instance, _ = provider({}, client=FakeClient(recall_texts=["Memory 1", "Memory 2"])) + instance.queue_prefetch("test") + if instance._prefetch_thread: + instance._prefetch_thread.join(timeout=5.0) + + result = instance.prefetch("test") + assert "Memory 1" in result + assert "Memory 2" in result + instance.shutdown() From 03bb459dfa8a5d44f2c4dc727459cfb05ba0de05 Mon Sep 17 00:00:00 2001 From: Yingliang Zhang Date: Wed, 23 Sep 2026 20:11:00 +0800 Subject: [PATCH 5/5] fix(hermes): bind retain work to enqueue-time session identity MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _make_turn_retain_job snapshotted turns/metadata/lineage/bank_id at enqueue but still built the aretain_batch item at run time via _build_retain_kwargs, which re-reads mutable per-session attributes (_retain_tags, _observation_scopes, _retain_source). A session switch landing between enqueue and writer drain stamped an OLD-session retain with the NEW session's tags/scopes/source — a document whose lineage tags contradicted its own snapshotted metadata session_id. Build the whole item at enqueue time; the writer job only ships it. Ported from NousResearch/hermes-agent#64499 (closed in the hindsight-move close pass). Tracked in #4662. --- hindsight-integrations/hermes/__init__.py | 23 +++++---- .../tests/test_retain_identity_isolation.py | 47 +++++++++++++++++++ 2 files changed, 61 insertions(+), 9 deletions(-) create mode 100644 hindsight-integrations/hermes/tests/test_retain_identity_isolation.py diff --git a/hindsight-integrations/hermes/__init__.py b/hindsight-integrations/hermes/__init__.py index d161191ea8..f6c6b0f6ed 100644 --- a/hindsight-integrations/hermes/__init__.py +++ b/hindsight-integrations/hermes/__init__.py @@ -1412,18 +1412,23 @@ def _retain_batch( def _make_turn_retain_job( self, turns: list[str], *, document_id: str, update_mode: str | None, label: str, track_ops: bool = True ) -> Callable[[], None]: - """Writer job shipping *turns* as one document. Inputs are snapshotted NOW: the - writer runs after later sync_turn() calls mutate _session_turns/_turn_index/_session_id.""" - content = "[" + ",".join(turns) + "]" - metadata = self._build_metadata(message_count=len(turns) * 2, turn_index=self._turn_index) + """Writer job shipping *turns* as one document. Everything is snapshotted NOW — + the item included — because the writer runs after later sync_turn() calls or an + on_session_switch() mutate the provider's mutable item-config attributes + (_retain_tags, _observation_scopes, _retain_source). Building the item at run + time would stamp an OLD-session retain with NEW-session tags/scopes.""" lineage = (("session", self._session_id), ("parent", self._parent_session_id)) tags = [f"{kind}:{sid}" for kind, sid in lineage if sid] or None - bank_id, retain_async, retain_context = self._bank_id, self._retain_async, self._retain_context + item = self._build_retain_kwargs( + "[" + ",".join(turns) + "]", + context=self._retain_context, + metadata=self._build_metadata(message_count=len(turns) * 2, turn_index=self._turn_index), + tags=tags, + update_mode=update_mode, + ) + bank_id, retain_async = self._bank_id, self._retain_async def _job() -> None: - item = self._build_retain_kwargs( - content, context=retain_context, metadata=metadata, tags=tags, update_mode=update_mode - ) logger.debug( "Hindsight %s: bank=%s, doc=%s, mode=%s, async=%s, content_len=%d, num_turns=%d", label, @@ -1431,7 +1436,7 @@ def _job() -> None: document_id, update_mode, retain_async, - len(content), + len(item["content"]), len(turns), ) resp = self._retain_batch(item, bank_id=bank_id, document_id=document_id, retain_async=retain_async) diff --git a/hindsight-integrations/hermes/tests/test_retain_identity_isolation.py b/hindsight-integrations/hermes/tests/test_retain_identity_isolation.py new file mode 100644 index 0000000000..1c55213298 --- /dev/null +++ b/hindsight-integrations/hermes/tests/test_retain_identity_isolation.py @@ -0,0 +1,47 @@ +"""Retain identity isolation — a queued retain must ship the identity and item +config that were authoritative at ENQUEUE time, not whatever the provider holds +when the writer thread eventually drains the queue +(NousResearch/hermes-agent#64499). + +``_make_turn_retain_job`` snapshots turns/metadata/lineage/bank_id at enqueue +time, but the writer job still calls ``_build_retain_kwargs`` at run time, which +re-reads mutable per-session attributes (``_retain_tags``, +``_observation_scopes``, ``_retain_source``). A session switch in between stamps +a queued OLD-session retain with the NEW session's item config.""" + +import threading + + +def test_queued_retain_keeps_enqueue_time_item_config(provider): + instance, fake = provider( + {"retain_tags": ["old-tag"], "retain_source": "old-source", "observation_scopes": "per_tag"} + ) + + # Park the writer exactly between dequeue and job execution so the main + # thread can mutate provider state while the job is still pending. + gate = threading.Event() + released = threading.Event() + + def _gate(): + gate.set() + released.wait(timeout=5.0) + + instance._retain_queue.put(_gate) + instance.sync_turn("old-user", "old-assistant") + gate.wait(timeout=5.0) + + # Mutate the item-config attributes the writer must NOT observe — what an + # on_session_switch() landing between enqueue and drain does. + instance._retain_tags = ["new-tag"] + instance._observation_scopes = "all_combinations" + instance._retain_source = "new-source" + + released.set() + instance._retain_queue.join() + + assert len(fake.retains) == 1 + item = fake.retains[0]["items"][0] + assert item["tags"] == ["old-tag", "session:session-1"] + assert item["observation_scopes"] == "per_tag" + assert item["metadata"]["source"] == "old-source" + instance.shutdown()