From f1273ed70461ff1c35d03f7f7004a5715d560e83 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Sat, 12 Sep 2026 23:21:47 +0530 Subject: [PATCH] refactor(honcho): one reclaim-key helper, one cache-size constant, no dead writer starter _SESSION_CACHE_MAX_SIZE was assigned twice with two comments describing one constant. _retain_for_retry and _keep_until_flushed shared the has-unsynced / current-owner / reclaim body and differed only in what to do when a newer object owns the key; _reclaim_key_locked returns that owner and each caller keeps its tail. _ensure_async_writer had no production caller (save() uses the _locked form under _async_thread_lock); removed, tests retargeted, constructor comment fixed. shutdown() sets _shutting_down under _async_thread_lock like the other site so save()'s flag check and the enqueue cannot interleave with it. tests/test_honcho_session_cache_bounds.py hoists its mid-file imports. --- plugins/memory/honcho/session.py | 53 ++++++++++------------- tests/honcho_plugin/test_async_memory.py | 16 +++---- tests/test_honcho_session_cache_bounds.py | 10 ++--- 3 files changed, 34 insertions(+), 45 deletions(-) diff --git a/plugins/memory/honcho/session.py b/plugins/memory/honcho/session.py index ffe5fa1394..4f82325278 100644 --- a/plugins/memory/honcho/session.py +++ b/plugins/memory/honcho/session.py @@ -25,13 +25,12 @@ logger = logging.getLogger(__name__) # Sentinel to signal the async writer thread to shut down _ASYNC_SHUTDOWN = object() -# Sessions remembered in _joined_author_peers; the oldest is dropped past this and its authors rejoin on their next write. -_SESSION_CACHE_MAX_SIZE = 128 # Honcho persists every message and get_or_create() re-hydrates on a miss, so the local copies stay bounded. _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. +# Hard caps for a burst of distinct sessions inside one TTL window; the dicts evict least recently used, +# and _joined_author_peers drops its oldest session so its authors rejoin on their next write. _SESSION_CACHE_MAX_SIZE = 128 _PEERS_CACHE_MAX_SIZE = 512 @@ -110,7 +109,7 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi self._prefetch_cache_lock = threading.Lock() # Async write queue — the writer thread starts lazily on first enqueue - # (_ensure_async_writer): constructing a manager must not spawn background + # (_ensure_async_writer_locked): constructing a manager must not spawn background # work or touch the network (unit tests build managers with mocked clients). self._async_queue: queue.Queue | None = queue.Queue() if self._write_frequency == "async" else None self._async_thread: threading.Thread | None = None @@ -494,31 +493,29 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi except Exception as e: logger.error("Honcho async writer error: %s", e) + def _reclaim_key_locked(self, session: HonchoSession) -> HonchoSession | None: + """Put a session with unsynced messages back where flush_all() looks. Returns the newer object that + owns the key when there is one (this session's batch still needs a home), None otherwise.""" + if not self._has_unsynced(session): + return None + current = self._cache.get(session.key) + if current is None: + self._cache[session.key] = session + return None if current is None or current is session else current + def _retain_for_retry(self, session: HonchoSession) -> None: - """Keep a session with unsynced messages where flush_all() looks. An evicted key takes it back; a key a - newer object owns keeps that object, and this one waits in the retry list.""" + """An evicted key takes the session back; a key a newer object owns keeps that object, and this one + waits in the retry list.""" with self._cache_lock: - if not self._has_unsynced(session): - return - current = self._cache.get(session.key) - if current is session or any(s is session for s in self._retry_sessions): - return - if current is None: - self._cache[session.key] = session - else: + if self._reclaim_key_locked(session) is not None and not any(s is session for s in self._retry_sessions): self._retry_sessions.append(session) def _keep_until_flushed(self, session: HonchoSession) -> None: - """Put an evicted session that still holds unsynced messages back where flush_all() looks. When a - newer object already owns the key, this one's batch is written now instead.""" + """Like _retain_for_retry, but when a newer object owns the key this batch is written now.""" with self._cache_lock: - current = self._cache.get(session.key) - if current is session or not self._has_unsynced(session): - return - if current is None: - self._cache[session.key] = session - return - self._flush_now(session) + owner = self._reclaim_key_locked(session) + if owner is not None: + self._flush_now(session) def _flush_now(self, session: HonchoSession) -> None: """A save-time flush whose failed batch stays reachable for flush_all(), even after an eviction.""" @@ -589,13 +586,6 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi skipped.append(item) return skipped - def _ensure_async_writer(self) -> None: - """Start the async writer on first enqueue (idempotent, thread-safe).""" - if self._async_thread is not None and self._async_thread.is_alive(): - return - with self._async_thread_lock: - self._ensure_async_writer_locked() - def _ensure_async_writer_locked(self) -> None: if self._async_thread is None or not self._async_thread.is_alive(): self._async_thread = spawn_context_thread(self._async_writer_loop, name="honcho-async-writer", owner=self) @@ -621,7 +611,8 @@ class HonchoSessionManager(SessionAuthMixin, SessionPeersMixin, SessionContextMi """Flush everything, then stop the async writer thread, within ``timeout``. The budget stops new uploads from starting and bounds the lock waits and the join; an upload already in flight runs to the client's HTTP timeout. Whatever stayed unsynced is counted in one warning.""" - self._shutting_down = True + with self._async_thread_lock: + self._shutting_down = True if self._async_queue is not None: deadline = time.monotonic() + timeout skipped = self._flush_cached_before(deadline) diff --git a/tests/honcho_plugin/test_async_memory.py b/tests/honcho_plugin/test_async_memory.py index 66013b8386..e779166573 100644 --- a/tests/honcho_plugin/test_async_memory.py +++ b/tests/honcho_plugin/test_async_memory.py @@ -289,14 +289,14 @@ class TestAsyncWriterThread: def test_shutdown_joins_thread(self, make_manager): mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() assert mgr._async_thread.is_alive() mgr.shutdown() assert not mgr._async_thread.is_alive() def test_async_writer_calls_flush(self, make_manager): mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() sess = _make_session() sess.add_message("user", "async msg") @@ -318,7 +318,7 @@ class TestAsyncWriterThread: def test_shutdown_sentinel_stops_loop(self, make_manager): mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() thread = mgr._async_thread mgr.shutdown() thread.join(timeout=10) @@ -331,7 +331,7 @@ class TestAsyncWriterThread: def test_stop_async_writer_joins_thread_without_flushing(self, make_manager): mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() sess = _make_session() sess.add_message("user", "must not be written") with mgr._cache_lock: @@ -432,7 +432,7 @@ class TestStopAsyncWriterDrain: class TestAsyncWriterRetry: def test_retries_once_on_failure(self, make_manager): mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() sess = _make_session() sess.add_message("user", "msg") @@ -459,7 +459,7 @@ class TestAsyncWriterRetry: """The shutdown flush already attempts the session within its budget; a 2s sleep and a second upload from the writer would run past it.""" mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() sess = _make_session() sess.add_message("user", "msg") calls = [] @@ -483,7 +483,7 @@ class TestAsyncWriterRetry: def test_drops_after_two_failures(self, make_manager): mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() sess = _make_session() sess.add_message("user", "msg") @@ -509,7 +509,7 @@ class TestAsyncWriterRetry: def test_retries_when_flush_reports_failure(self, make_manager): mgr = make_manager(write_frequency="async") - mgr._ensure_async_writer() + mgr._ensure_async_writer_locked() sess = _make_session() sess.add_message("user", "msg") diff --git a/tests/test_honcho_session_cache_bounds.py b/tests/test_honcho_session_cache_bounds.py index d54102ae87..e1b35f3b1f 100644 --- a/tests/test_honcho_session_cache_bounds.py +++ b/tests/test_honcho_session_cache_bounds.py @@ -11,9 +11,13 @@ import time from datetime import datetime, timedelta from types import SimpleNamespace +import pytest + from plugins.memory.honcho.session import ( HonchoSession, HonchoSessionManager, + _PEERS_CACHE_MAX_SIZE, + _SESSION_CACHE_MAX_SIZE, _SESSION_IDLE_TTL_SECONDS, _SESSION_MESSAGE_RETENTION, ) @@ -161,12 +165,6 @@ def test_get_or_create_triggers_sweep_without_blocking_on_lock_reentrancy(): # hard caps, unsynced buffers, peers, and read activity (follows #71463) # --------------------------------------------------------------------------- -import pytest # noqa: E402 - -from plugins.memory.honcho.session import ( # noqa: E402 - _PEERS_CACHE_MAX_SIZE, - _SESSION_CACHE_MAX_SIZE, -) def _fill_sessions(mgr, count, unsynced_keys=()):