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.
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
|
||||
|
||||
@@ -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=()):
|
||||
|
||||
Reference in New Issue
Block a user