Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
a98252f
fix(hermes): guard client lifecycle with a leaf lock
yingliang-zhang Sep 23, 2026
3241d8e
chore: sync generated files for verify-generated-files (mirrors #4678)
yingliang-zhang Sep 23, 2026
3616b6c
Merge origin/main into fix/hindsight-client-leaf-lock
yingliang-zhang Sep 30, 2026
ac0d77b
fix(hermes): pass retain_async=True to aretain_batch in tool handler
yingliang-zhang Sep 23, 2026
5a5096f
fix(hermes): fence prefetch publication to the owning session generation
yingliang-zhang Sep 23, 2026
03bb459
fix(hermes): bind retain work to enqueue-time session identity
yingliang-zhang Sep 23, 2026
803585b
fix(hermes): gate prefetch spawns on the session that owns the slot
yingliang-zhang Sep 23, 2026
4664a6d
fix(hermes): frame recalled context as untrusted data
yingliang-zhang Sep 23, 2026
aba3045
chore: collapse _MALICIOUS_RECALL into a single string literal (mirro…
yingliang-zhang Sep 23, 2026
b6721a3
Merge remote-tracking branch 'origin/main' into fix/hindsight-client-…
yingliang-zhang Sep 30, 2026
83f4104
Merge branch 'fix/hindsight-client-leaf-lock' into fix/hindsight-reta…
yingliang-zhang Sep 30, 2026
d6742eb
Merge branch 'fix/hindsight-retain-async-tool-handler' into fix/hinds…
yingliang-zhang Sep 30, 2026
37735ec
Merge branch 'fix/hindsight-prefetch-session-identity' into fix/hinds…
yingliang-zhang Sep 30, 2026
4527350
Merge branch 'fix/hindsight-retain-session-identity' into fix/hindsig…
yingliang-zhang Sep 30, 2026
a07b32a
Merge branch 'fix/hindsight-prefetch-session-owner' into fix/memory-u…
yingliang-zhang Sep 30, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
151 changes: 125 additions & 26 deletions hindsight-integrations/hermes/__init__.py

Large diffs are not rendered by default.

49 changes: 48 additions & 1 deletion hindsight-integrations/hermes/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
import json
import logging
import re
from typing import Any, List
from typing import Any, Dict, List

# Log under the plugin package's own logger name (loader-path independent).
logger = logging.getLogger(__name__.rpartition(".")[0])
Expand All @@ -30,6 +30,8 @@
# (vectorize-io/hindsight#932).
_MIN_VERSION_FOR_UPDATE_MODE_APPEND = "0.5.0"
_VALID_BUDGETS = {"low", "mid", "high"}
_PREFETCH_JSON_CHARS_PER_TOKEN = 4
_MIN_PREFETCH_JSON_CHARS = 256
_PROVIDER_DEFAULT_MODELS = {
"openai": "gpt-4o-mini",
"anthropic": "claude-haiku-4-5",
Expand Down Expand Up @@ -57,6 +59,51 @@ def _parse_int_setting(value: Any, default: int) -> int:
return default


def _serialize_prefetch_data(kind: str, content: List[str], *, max_chars: int) -> str:
"""Serialize untrusted Hindsight text as bounded JSON reference data."""
payload: Dict[str, Any] = {
"source": "hindsight",
"kind": kind,
"content": [],
}
limit = max(_MIN_PREFETCH_JSON_CHARS, int(max_chars))

def _encode() -> str:
encoded = json.dumps(payload, ensure_ascii=False, separators=(",", ":"))
return encoded.replace("<", "\\u003c").replace(">", "\\u003e")

for raw_text in content:
if not raw_text:
continue
text = str(raw_text)
payload["content"].append(text)
if len(_encode()) <= limit:
continue
payload["content"].pop()

low = 0
high = len(text)
best = ""
while low <= high:
midpoint = (low + high) // 2
excerpt = text[:midpoint]
if midpoint < len(text):
excerpt += "…"
payload["content"].append(excerpt)
candidate = _encode()
payload["content"].pop()
if len(candidate) <= limit:
best = excerpt
low = midpoint + 1
else:
high = midpoint - 1
if best:
payload["content"].append(best)
break

return _encode() if payload["content"] else ""


def _daemon_llm_provider(provider: str) -> str:
return "openai" if provider in _OPENAI_WIRE_PROVIDERS else provider

Expand Down
159 changes: 159 additions & 0 deletions hindsight-integrations/hermes/tests/test_client_lifecycle.py
Original file line number Diff line number Diff line change
@@ -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
Loading
Loading