diff --git a/plugins/memory/honcho/session.py b/plugins/memory/honcho/session.py index e64916c18f..75c50398ba 100644 --- a/plugins/memory/honcho/session.py +++ b/plugins/memory/honcho/session.py @@ -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: diff --git a/plugins/memory/honcho/session_context.py b/plugins/memory/honcho/session_context.py index 824cc83442..c39fb42d39 100644 --- a/plugins/memory/honcho/session_context.py +++ b/plugins/memory/honcho/session_context.py @@ -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 "" diff --git a/plugins/memory/honcho/session_migration.py b/plugins/memory/honcho/session_migration.py index 575dfb8639..c706107a2d 100644 --- a/plugins/memory/honcho/session_migration.py +++ b/plugins/memory/honcho/session_migration.py @@ -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 diff --git a/tests/honcho_plugin/test_pin_peer_name.py b/tests/honcho_plugin/test_pin_peer_name.py index 4e481ba324..7530aba74c 100644 --- a/tests/honcho_plugin/test_pin_peer_name.py +++ b/tests/honcho_plugin/test_pin_peer_name.py @@ -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) ) diff --git a/tests/honcho_plugin/test_session.py b/tests/honcho_plugin/test_session.py index 95b281230c..16764f6c26 100644 --- a/tests/honcho_plugin/test_session.py +++ b/tests/honcho_plugin/test_session.py @@ -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 = {} diff --git a/tests/test_honcho_session_cache_bounds.py b/tests/test_honcho_session_cache_bounds.py index 46c6ec3b8b..54be4ac7ea 100644 --- a/tests/test_honcho_session_cache_bounds.py +++ b/tests/test_honcho_session_cache_bounds.py @@ -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)