fix(honcho): cap the manager caches and keep unsynced sessions out of eviction
the idle sweep from #71463 left three growth paths open. _peers_cache had no bound at all, _session_observation (from #98941) grew one entry per session id and kept orphans when an init failed after add_peers, and a burst of distinct sessions inside one ttl window was not bounded. the sweep also evicted sessions whose messages had not reached honcho yet, which in "session" write mode drops the only copy, and _flush_session re-inserted an evicted session into the cache with no observation flags, so recall for it routed from the config snapshot instead of its server config. _cache and _sessions_cache now cap at 128 entries and _peers_cache at 512, evicting least recently used first (dict order, refreshed on every hit). a session with unsynced messages is never evicted by the sweep or the cap. _configure_session_peers returns the synced flags and get_or_create stores them under _cache_lock next to the cache entry, so the observation dict can hold no id the cache does not; eviction drops both. recall reads go through _cached_session, which stamps updated_at, so a read-only session survives the idle sweep. follows #71461, #71463, #98936.
This commit is contained in:
@@ -33,6 +33,9 @@ _SESSION_CACHE_MAX_SIZE = 128
|
||||
_SESSION_MESSAGE_RETENTION = 200
|
||||
_SESSION_IDLE_TTL_SECONDS = 3600
|
||||
_SESSION_SWEEP_INTERVAL_SECONDS = 300
|
||||
# Hard caps for a burst of distinct sessions inside one TTL window; both dicts evict least recently used.
|
||||
_SESSION_CACHE_MAX_SIZE = 128
|
||||
_PEERS_CACHE_MAX_SIZE = 512
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -133,6 +136,7 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
while True:
|
||||
with self._cache_lock:
|
||||
if key in cache:
|
||||
cache[key] = cache.pop(key) # dict order is the LRU order the caps evict from
|
||||
return cache[key]
|
||||
generation = self._client_generation
|
||||
obj = fetch()
|
||||
@@ -141,6 +145,15 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
return cache.setdefault(key, obj)
|
||||
# Client rebuilt mid-resolve: this object holds the discarded transport. Retry.
|
||||
|
||||
def _cached_session(self, session_key: str) -> HonchoSession | None:
|
||||
"""The locally cached session, stamped as used so recall alone keeps it out of the idle sweep."""
|
||||
with self._cache_lock:
|
||||
session = self._cache.get(session_key)
|
||||
if session is not None:
|
||||
session.updated_at = datetime.now()
|
||||
self._cache[session_key] = self._cache.pop(session_key)
|
||||
return session
|
||||
|
||||
def _sdk_session(self, session_id: str) -> Any:
|
||||
"""Get or create the SDK session (cached until a client rebuild clears the cache)."""
|
||||
return self._cached_sdk_object(self._sessions_cache, session_id, lambda: self.honcho.session(session_id))
|
||||
@@ -168,12 +181,14 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
|
||||
# ----- Session creation -----
|
||||
|
||||
def _configure_session_peers(self, session_id: str, user_peer: Any, assistant_peer: Any) -> bool:
|
||||
def _configure_session_peers(self, session_id: str, user_peer: Any, assistant_peer: Any) -> dict[str, bool] | None:
|
||||
"""add_peers with this session's observation config, then adopt the server's effective
|
||||
config (set via the Honcho UI, it wins over local defaults) under the session's own id.
|
||||
Returns False when auth died mid-way (already recorded by _authed_call)."""
|
||||
config (set via the Honcho UI, it wins over local defaults). Returns the effective flags for
|
||||
get_or_create to store next to the cache entry, or None when auth died mid-way (already
|
||||
recorded by _authed_call)."""
|
||||
peers = (("user", user_peer), ("ai", assistant_peer))
|
||||
flags = self._observation_flags(session_id)
|
||||
synced = dict(flags)
|
||||
try:
|
||||
from honcho.session import SessionPeerConfig
|
||||
peer_entries = [
|
||||
@@ -187,13 +202,11 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
"peer configuration read",
|
||||
lambda: [self._sdk_session(session_id).get_peer_configuration(peer) for _, peer in peers],
|
||||
)
|
||||
synced = dict(flags)
|
||||
for (kind, _), server_cfg in zip(peers, server_cfgs):
|
||||
for field_name in ("observe_me", "observe_others"):
|
||||
value = getattr(server_cfg, field_name)
|
||||
if value is not None:
|
||||
synced[f"{kind}_{field_name}"] = value
|
||||
self._session_observation[session_id] = synced
|
||||
logger.debug("Honcho observation synced from server for session '%s': user(me=%s,others=%s) ai(me=%s,others=%s)",
|
||||
session_id, synced["user_observe_me"], synced["user_observe_others"],
|
||||
synced["ai_observe_me"], synced["ai_observe_others"])
|
||||
@@ -201,10 +214,10 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
self._guarded(_adopt_server_config, None, logging.DEBUG,
|
||||
"Honcho get_peer_configuration failed (using local config): %s")
|
||||
except HonchoAuthError:
|
||||
return False
|
||||
return None
|
||||
except Exception as e:
|
||||
logger.warning("Honcho session '%s' add_peers failed (non-fatal): %s", session_id, e)
|
||||
return True
|
||||
return synced
|
||||
|
||||
def _load_existing_messages(self, session_id: str) -> list:
|
||||
"""Load prior messages via context() (one call for messages + metadata), oldest first."""
|
||||
@@ -229,36 +242,74 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
logger.warning("Honcho session '%s' loaded (failed to fetch context: %s)", session_id, e)
|
||||
return []
|
||||
|
||||
def _get_or_create_honcho_session(self, session_id: str, user_peer: Any, assistant_peer: Any) -> tuple[Any, list]:
|
||||
"""(honcho_session, existing_messages) with peers configured; a cached session yields no messages."""
|
||||
def _get_or_create_honcho_session(
|
||||
self, session_id: str, user_peer: Any, assistant_peer: Any,
|
||||
) -> tuple[Any, list, dict[str, bool] | None]:
|
||||
"""(honcho_session, existing_messages, observation flags) with peers configured; a cached
|
||||
session yields no messages and no flags."""
|
||||
with self._cache_lock:
|
||||
if session_id in self._sessions_cache:
|
||||
logger.debug("Honcho session '%s' retrieved from cache", session_id)
|
||||
return self._sessions_cache[session_id], []
|
||||
return self._sessions_cache[session_id], [], None
|
||||
|
||||
self._authed_call("session setup", lambda: self._sdk_session(session_id))
|
||||
existing_messages: list = (self._load_existing_messages(session_id)
|
||||
if self._configure_session_peers(session_id, user_peer, assistant_peer) else [])
|
||||
observation = self._configure_session_peers(session_id, user_peer, assistant_peer)
|
||||
existing_messages: list = self._load_existing_messages(session_id) if observation is not None else []
|
||||
|
||||
with self._cache_lock:
|
||||
honcho_session = self._sessions_cache.get(session_id)
|
||||
if honcho_session is None:
|
||||
# A mid-init client rebuild dropped the cached session; resolve a fresh one.
|
||||
honcho_session = self._authed_call("session setup", lambda: self._sdk_session(session_id))
|
||||
return honcho_session, existing_messages
|
||||
return honcho_session, existing_messages, observation
|
||||
|
||||
@staticmethod
|
||||
def _has_unsynced(session: HonchoSession) -> bool:
|
||||
return any(not m.get("_synced") for m in session.messages)
|
||||
|
||||
def _evict_session_locked(self, key: str, session: HonchoSession) -> None:
|
||||
"""Drop one session and every entry keyed to it. Caller holds _cache_lock."""
|
||||
del self._cache[key]
|
||||
self._sessions_cache.pop(session.honcho_session_id, None)
|
||||
self._session_observation.pop(session.honcho_session_id, None)
|
||||
with self._prefetch_cache_lock:
|
||||
self._context_cache.pop(key, None)
|
||||
|
||||
def _enforce_cache_caps_locked(self) -> None:
|
||||
"""Evict least recently used entries above the hard caps. A session with unsynced messages
|
||||
is never evicted: it is the only copy until the flush lands. Caller holds _cache_lock."""
|
||||
for key, session in list(self._cache.items()):
|
||||
if len(self._cache) <= _SESSION_CACHE_MAX_SIZE:
|
||||
break
|
||||
if not self._has_unsynced(session):
|
||||
self._evict_session_locked(key, session)
|
||||
live_ids = {s.honcho_session_id for s in self._cache.values()}
|
||||
for session_id in list(self._sessions_cache):
|
||||
if len(self._sessions_cache) <= _SESSION_CACHE_MAX_SIZE:
|
||||
break
|
||||
if session_id not in live_ids:
|
||||
del self._sessions_cache[session_id]
|
||||
live_peers = {p for s in self._cache.values() for p in (s.user_peer_id, s.assistant_peer_id)}
|
||||
for peer_id in list(self._peers_cache):
|
||||
if len(self._peers_cache) <= _PEERS_CACHE_MAX_SIZE:
|
||||
break
|
||||
if peer_id not in live_peers:
|
||||
del self._peers_cache[peer_id]
|
||||
with self._prefetch_cache_lock:
|
||||
for key in [k for k in self._context_cache if k not in self._cache]:
|
||||
del self._context_cache[key]
|
||||
|
||||
def _sweep_idle_sessions_locked(self) -> int:
|
||||
"""Evict sessions idle beyond _SESSION_IDLE_TTL_SECONDS together with their SDK session and
|
||||
prefetch entries. Caller holds _cache_lock."""
|
||||
"""Evict sessions idle beyond _SESSION_IDLE_TTL_SECONDS together with their SDK session,
|
||||
observation flags and prefetch entries, then enforce the caps. Caller holds _cache_lock."""
|
||||
cutoff = time.time() - _SESSION_IDLE_TTL_SECONDS
|
||||
evicted = 0
|
||||
for key, session in list(self._cache.items()):
|
||||
if session.updated_at.timestamp() >= cutoff:
|
||||
if session.updated_at.timestamp() >= cutoff or self._has_unsynced(session):
|
||||
continue
|
||||
del self._cache[key]
|
||||
self._sessions_cache.pop(session.honcho_session_id, None)
|
||||
self._context_cache.pop(key, None)
|
||||
self._evict_session_locked(key, session)
|
||||
evicted += 1
|
||||
self._enforce_cache_caps_locked()
|
||||
return evicted
|
||||
|
||||
def _maybe_sweep_idle_sessions(self) -> None:
|
||||
@@ -277,10 +328,9 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
``user_peer_id`` replaces the resolved user peer when the session's participant is not the
|
||||
runtime user, e.g. the sender bot of an a2a session."""
|
||||
self._maybe_sweep_idle_sessions()
|
||||
with self._cache_lock:
|
||||
if key in self._cache:
|
||||
logger.debug("Local session cache hit: %s", key)
|
||||
return self._cache[key]
|
||||
if (cached := self._cached_session(key)) is not None:
|
||||
logger.debug("Local session cache hit: %s", key)
|
||||
return cached
|
||||
|
||||
# Gateway sessions normally use the platform-native runtime identity so multi-user
|
||||
# bots scope memory per user; config can alias/prefix it, or pinPeerName pins all
|
||||
@@ -293,7 +343,7 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
honcho_session_id = self._sanitize_id(key)
|
||||
user_peer = self._get_or_create_peer(user_peer_id)
|
||||
assistant_peer = self._get_or_create_peer(assistant_peer_id)
|
||||
_, existing_messages = self._get_or_create_honcho_session(honcho_session_id, user_peer, assistant_peer)
|
||||
_, existing_messages, observation = self._get_or_create_honcho_session(honcho_session_id, user_peer, assistant_peer)
|
||||
|
||||
session = HonchoSession(
|
||||
key=key, user_peer_id=user_peer_id, assistant_peer_id=assistant_peer_id, honcho_session_id=honcho_session_id,
|
||||
@@ -305,6 +355,10 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
)
|
||||
with self._cache_lock:
|
||||
self._cache[key] = session
|
||||
# Stored with the cache entry and dropped with it, so the dict can hold no orphan ids.
|
||||
if observation is not None:
|
||||
self._session_observation[honcho_session_id] = observation
|
||||
self._enforce_cache_caps_locked()
|
||||
return session
|
||||
|
||||
# ----- Writes -----
|
||||
@@ -363,7 +417,7 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
assistant_peer = self._get_or_create_peer(session.assistant_peer_id)
|
||||
honcho_session = self._sessions_cache.get(session.honcho_session_id)
|
||||
if honcho_session is None:
|
||||
honcho_session, _ = self._get_or_create_honcho_session(session.honcho_session_id, user_peer, assistant_peer)
|
||||
honcho_session, _, _ = self._get_or_create_honcho_session(session.honcho_session_id, user_peer, assistant_peer)
|
||||
honcho_messages = []
|
||||
for m in new_messages:
|
||||
if m["role"] != "user":
|
||||
@@ -386,8 +440,6 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi
|
||||
msg["_synced"] = ok
|
||||
if ok:
|
||||
self._trim_synced_messages(session)
|
||||
with self._cache_lock:
|
||||
self._cache[session.key] = session
|
||||
return ok
|
||||
|
||||
def _try_flush(self, session: HonchoSession, level: int, msg: str) -> bool:
|
||||
|
||||
@@ -64,7 +64,7 @@ class SessionContextMixin:
|
||||
self, session_key: str, fn: Callable[[Any], Any], default: Any, level: int, msg: str, *args: Any,
|
||||
) -> Any:
|
||||
"""``_guarded`` over ``fn(session)`` for the cached session; ``default`` when no session is cached."""
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
return self._guarded(lambda: fn(session), default, level, msg, *args) if session else default
|
||||
|
||||
@staticmethod
|
||||
@@ -126,7 +126,7 @@ class SessionContextMixin:
|
||||
"""Pre-fetch user + AI peer context (representation, card) plus the session summary.
|
||||
``user_message`` is passed as search_query so Honcho returns topic-relevant conclusions.
|
||||
Stops early (returning what it has) once auth is dead."""
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
if not session:
|
||||
return {}
|
||||
result: dict[str, str] = {}
|
||||
@@ -180,7 +180,7 @@ class SessionContextMixin:
|
||||
def get_session_context(self, session_key: str, peer: str = "user") -> dict[str, Any]:
|
||||
"""Fetch session-level context (summary, representation, card, recent messages).
|
||||
Raises HonchoAuthError so callers can tell rejected credentials from no context."""
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
if not session:
|
||||
return {}
|
||||
if session.honcho_session_id not in self._sessions_cache:
|
||||
@@ -225,7 +225,7 @@ class SessionContextMixin:
|
||||
"""Hybrid search over raw messages visible from ``peer``'s perspective, all sessions. Snippets
|
||||
accumulate until ``max_tokens`` (~4 chars/token) is exhausted. Returns "" when nothing matches;
|
||||
raises HonchoAuthError on rejected credentials."""
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
q = (query or "").strip()[:4000] # Honcho caps query length for the embedding model.
|
||||
if not session or not q:
|
||||
return ""
|
||||
@@ -277,7 +277,7 @@ class SessionContextMixin:
|
||||
"""Write a conclusion (durable fact) about ``peer`` back to Honcho."""
|
||||
if not content or not content.strip():
|
||||
return False
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
if not session:
|
||||
logger.warning("No session cached for '%s', skipping conclusion", session_key)
|
||||
return False
|
||||
@@ -341,7 +341,7 @@ class SessionContextMixin:
|
||||
operations, auth failures are logged and swallowed here too."""
|
||||
if not content or not content.strip():
|
||||
return False
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
if not session:
|
||||
logger.warning("No session cached for '%s', skipping AI seed", session_key)
|
||||
return False
|
||||
@@ -390,7 +390,7 @@ class SessionContextMixin:
|
||||
server error surfaces as an error, not as "no result" (#36098 issue 4: collapsing failures to ""
|
||||
made auth errors, timeouts, and genuinely-empty answers indistinguishable).
|
||||
"""
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
target_peer_id = self._resolve_peer_id(session, peer) if session else None
|
||||
if target_peer_id is None:
|
||||
return ""
|
||||
|
||||
@@ -25,7 +25,7 @@ class SessionMigrationMixin:
|
||||
if not memory_path.exists():
|
||||
return False
|
||||
|
||||
session = self._cache.get(session_key)
|
||||
session = self._cached_session(session_key)
|
||||
if not session:
|
||||
logger.warning("No local session cached for '%s', skipping memory migration", session_key)
|
||||
return False
|
||||
|
||||
@@ -108,7 +108,7 @@ def _patch_manager_for_resolution_test(mgr: HonchoSessionManager) -> None:
|
||||
fake_peer = MagicMock()
|
||||
mgr._get_or_create_peer = MagicMock(return_value=fake_peer)
|
||||
mgr._get_or_create_honcho_session = MagicMock(
|
||||
return_value=(MagicMock(), [])
|
||||
return_value=(MagicMock(), [], None)
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -1164,6 +1164,7 @@ class TestSetPeerCardNoneGuard:
|
||||
cfg = HonchoClientConfig(api_key="test-key", enabled=True)
|
||||
mgr = HonchoSessionManager.__new__(HonchoSessionManager)
|
||||
mgr._cache = {}
|
||||
mgr._cache_lock = threading.RLock()
|
||||
mgr._sessions_cache = {}
|
||||
mgr._config = cfg
|
||||
mgr._session_observation = {}
|
||||
@@ -1203,6 +1204,7 @@ class TestGetSessionContextFallback:
|
||||
cfg = HonchoClientConfig(api_key="test-key", enabled=True)
|
||||
mgr = HonchoSessionManager.__new__(HonchoSessionManager)
|
||||
mgr._cache = {}
|
||||
mgr._cache_lock = threading.RLock()
|
||||
mgr._sessions_cache = {}
|
||||
mgr._config = cfg
|
||||
mgr._dialectic_dynamic = True
|
||||
@@ -1520,22 +1522,22 @@ class TestObservationPerSessionScoping:
|
||||
)
|
||||
|
||||
def test_sync_back_scopes_flags_per_session(self):
|
||||
"""Server flags land under each session's own id; manager snapshot untouched."""
|
||||
"""Server flags come back per session; manager snapshot untouched."""
|
||||
mgr = self._make_manager()
|
||||
|
||||
# Session A's server config disables user observe_others...
|
||||
self._setup_session(
|
||||
_, _, flags_a = self._setup_session(
|
||||
mgr, "sid-a",
|
||||
_FakeSdkSession(_FakeServerPeerConfig(observe_others=False), _FakeServerPeerConfig()),
|
||||
)
|
||||
# ...session B's server leaves everything at the synced-in defaults.
|
||||
self._setup_session(
|
||||
_, _, flags_b = self._setup_session(
|
||||
mgr, "sid-b",
|
||||
_FakeSdkSession(_FakeServerPeerConfig(), _FakeServerPeerConfig()),
|
||||
)
|
||||
|
||||
assert mgr._session_observation["sid-a"]["user_observe_others"] is False
|
||||
assert mgr._session_observation["sid-b"]["user_observe_others"] is True
|
||||
assert flags_a["user_observe_others"] is False
|
||||
assert flags_b["user_observe_others"] is True
|
||||
# The config snapshot on the manager must survive both syncs — this is
|
||||
# the regression: last-session-wins used to overwrite it (#98936).
|
||||
assert mgr._user_observe_others is True
|
||||
@@ -1547,7 +1549,8 @@ class TestObservationPerSessionScoping:
|
||||
fake = _FakeSdkSession(
|
||||
_FakeServerPeerConfig(observe_others=False), _FakeServerPeerConfig()
|
||||
)
|
||||
self._setup_session(mgr, "sid-a", fake)
|
||||
_, _, flags = self._setup_session(mgr, "sid-a", fake)
|
||||
mgr._session_observation["sid-a"] = flags # what get_or_create stores next to the cache entry
|
||||
|
||||
# Force the full setup path again (cache cleared, e.g. after re-auth).
|
||||
mgr._sessions_cache = {}
|
||||
|
||||
@@ -155,3 +155,138 @@ def test_get_or_create_triggers_sweep_without_blocking_on_lock_reentrancy():
|
||||
t.join(timeout=5)
|
||||
assert done.is_set(), "sweep did not complete — possible deadlock"
|
||||
assert "stale" not in mgr._cache
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# hard caps, unsynced buffers, peers, and read activity (follows #71463)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
from plugins.memory.honcho.session import ( # noqa: E402
|
||||
_PEERS_CACHE_MAX_SIZE,
|
||||
_SESSION_CACHE_MAX_SIZE,
|
||||
)
|
||||
|
||||
|
||||
def _fill_sessions(mgr, count, unsynced_keys=()):
|
||||
for i in range(count):
|
||||
key = f"k{i}"
|
||||
session = _session(key=key)
|
||||
if key in unsynced_keys:
|
||||
session.add_message("user", "pending", _synced=False)
|
||||
mgr._cache[key] = session
|
||||
mgr._sessions_cache[session.honcho_session_id] = object()
|
||||
mgr._session_observation[session.honcho_session_id] = {"ai_observe_others": False}
|
||||
mgr._context_cache[key] = {"representation": "r"}
|
||||
|
||||
|
||||
def test_size_cap_evicts_least_recently_used_sessions_with_their_entries():
|
||||
mgr = _manager()
|
||||
_fill_sessions(mgr, _SESSION_CACHE_MAX_SIZE + 2)
|
||||
|
||||
with mgr._cache_lock:
|
||||
mgr._enforce_cache_caps_locked()
|
||||
|
||||
assert len(mgr._cache) == _SESSION_CACHE_MAX_SIZE
|
||||
assert "k0" not in mgr._cache and "k1" not in mgr._cache
|
||||
assert "k2" in mgr._cache
|
||||
for gone in ("k0", "k1"):
|
||||
assert f"hs-{gone}" not in mgr._sessions_cache
|
||||
assert f"hs-{gone}" not in mgr._session_observation
|
||||
assert gone not in mgr._context_cache
|
||||
assert "hs-k2" in mgr._sessions_cache and "hs-k2" in mgr._session_observation
|
||||
|
||||
|
||||
def test_size_cap_never_evicts_a_session_with_unsynced_messages():
|
||||
mgr = _manager()
|
||||
_fill_sessions(mgr, _SESSION_CACHE_MAX_SIZE + 1, unsynced_keys={"k0"})
|
||||
|
||||
with mgr._cache_lock:
|
||||
mgr._enforce_cache_caps_locked()
|
||||
|
||||
assert "k0" in mgr._cache # the only copy until its flush lands
|
||||
assert "k1" not in mgr._cache
|
||||
assert len(mgr._cache) == _SESSION_CACHE_MAX_SIZE
|
||||
|
||||
|
||||
def test_idle_sweep_keeps_sessions_with_unsynced_messages():
|
||||
mgr = _manager()
|
||||
stale = _session(key="stale")
|
||||
stale.add_message("user", "pending", _synced=False)
|
||||
stale.updated_at = datetime.now() - timedelta(seconds=_SESSION_IDLE_TTL_SECONDS + 1)
|
||||
mgr._cache = {"stale": stale}
|
||||
|
||||
with mgr._cache_lock:
|
||||
evicted = mgr._sweep_idle_sessions_locked()
|
||||
|
||||
assert evicted == 0
|
||||
assert "stale" in mgr._cache
|
||||
|
||||
|
||||
def test_peers_cap_evicts_unreferenced_peers_oldest_first():
|
||||
mgr = _manager()
|
||||
live = _session(key="live")
|
||||
mgr._cache = {"live": live}
|
||||
mgr._peers_cache[live.user_peer_id] = object()
|
||||
mgr._peers_cache[live.assistant_peer_id] = object()
|
||||
for i in range(_PEERS_CACHE_MAX_SIZE):
|
||||
mgr._peers_cache[f"guest{i}"] = object()
|
||||
|
||||
with mgr._cache_lock:
|
||||
mgr._enforce_cache_caps_locked()
|
||||
|
||||
assert len(mgr._peers_cache) == _PEERS_CACHE_MAX_SIZE
|
||||
assert live.user_peer_id in mgr._peers_cache and live.assistant_peer_id in mgr._peers_cache
|
||||
assert "guest0" not in mgr._peers_cache and "guest1" not in mgr._peers_cache
|
||||
assert f"guest{_PEERS_CACHE_MAX_SIZE - 1}" in mgr._peers_cache
|
||||
|
||||
|
||||
def test_recall_read_counts_as_activity_for_the_idle_sweep():
|
||||
mgr = _manager()
|
||||
session = _session(key="read-only")
|
||||
session.updated_at = datetime.now() - timedelta(seconds=_SESSION_IDLE_TTL_SECONDS + 1)
|
||||
mgr._cache = {"read-only": session}
|
||||
|
||||
assert mgr._cached_session("read-only") is session
|
||||
with mgr._cache_lock:
|
||||
evicted = mgr._sweep_idle_sessions_locked()
|
||||
|
||||
assert evicted == 0
|
||||
assert "read-only" in mgr._cache
|
||||
|
||||
|
||||
def test_sdk_object_hit_moves_the_key_to_the_recent_end():
|
||||
mgr = _manager()
|
||||
mgr._peers_cache = {"a": object(), "b": object()}
|
||||
|
||||
mgr._cached_sdk_object(mgr._peers_cache, "a", lambda: None)
|
||||
|
||||
assert list(mgr._peers_cache) == ["b", "a"]
|
||||
|
||||
|
||||
def test_get_or_create_stores_observation_flags_with_the_entry_and_eviction_drops_them():
|
||||
mgr = _manager()
|
||||
mgr._config.ai_peer = "hermes"
|
||||
flags = {"user_observe_me": True, "user_observe_others": True, "ai_observe_me": True, "ai_observe_others": False}
|
||||
mgr._get_or_create_peer = lambda peer_id: object()
|
||||
mgr._get_or_create_honcho_session = lambda sid, user, assistant: (object(), [], dict(flags))
|
||||
|
||||
session = mgr.get_or_create("cli:one")
|
||||
|
||||
assert mgr._session_observation[session.honcho_session_id] == flags
|
||||
assert mgr._ai_observes_others(session) is False
|
||||
with mgr._cache_lock:
|
||||
mgr._evict_session_locked("cli:one", session)
|
||||
assert session.honcho_session_id not in mgr._session_observation
|
||||
|
||||
|
||||
def test_flush_does_not_resurrect_an_evicted_session():
|
||||
mgr = _manager()
|
||||
session = _session(key="gone")
|
||||
session.add_message("user", "late", _synced=False)
|
||||
peer = SimpleNamespace(message=lambda content: content)
|
||||
mgr._get_or_create_peer = lambda peer_id: peer
|
||||
mgr._sessions_cache[session.honcho_session_id] = SimpleNamespace(add_messages=lambda messages: None)
|
||||
|
||||
assert mgr._flush_session(session) is True
|
||||
assert "gone" not in mgr._cache
|
||||
assert all(m["_synced"] for m in session.messages)
|
||||
|
||||
Reference in New Issue
Block a user