fix(honcho): join the session manager's async-writer thread on provider shutdown
Provider shutdown() only called manager.flush_all(), which drains the queue but never joins the async-writer thread — manager.shutdown() exists and nothing called it. The writer thread could still be blocked in httpx I/O at interpreter exit (the #37632 crash class). Now shutdown() calls manager.shutdown() (flush + join) when persistence is enabled, and a new manager.stop_async_writer() (join only, no flush) when saveMessages is false, so containment and clean teardown compose.
This commit is contained in:
@@ -1672,15 +1672,26 @@ class HonchoMemoryProvider(MemoryProvider):
|
||||
for t in (self._prefetch_thread, self._sync_thread):
|
||||
if t and t.is_alive():
|
||||
t.join(timeout=5.0)
|
||||
# Flush any remaining messages. Honors saveMessages: false — skip
|
||||
# persistence, but the worker-thread joins above still run (cleanup
|
||||
# is independent of persistence; placing the guard here rather than
|
||||
# at the top avoids leaking _prefetch_thread/_sync_thread).
|
||||
manager = self._manager
|
||||
if manager and self._init_thread and self._init_thread.is_alive() and not self._session_initialized:
|
||||
manager = None
|
||||
# Honors saveMessages: false — skip persistence, but thread cleanup
|
||||
# still runs: the session manager's async-writer thread must be
|
||||
# joined either way so daemon threads aren't left blocked in httpx
|
||||
# I/O during interpreter finalization.
|
||||
if not getattr(self._config, "save_messages", True):
|
||||
if manager:
|
||||
try:
|
||||
manager.stop_async_writer()
|
||||
except Exception:
|
||||
pass
|
||||
return
|
||||
if self._manager and not (self._init_thread and self._init_thread.is_alive() and not self._session_initialized):
|
||||
if manager:
|
||||
try:
|
||||
self._manager.flush_all()
|
||||
# manager.shutdown() = flush_all() + join the async-writer
|
||||
# thread. Previously only flush_all() ran here, leaving the
|
||||
# writer thread alive at exit.
|
||||
manager.shutdown()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@@ -785,6 +785,18 @@ class HonchoSessionManager:
|
||||
)
|
||||
self._async_thread.start()
|
||||
|
||||
def stop_async_writer(self) -> None:
|
||||
"""Stop the async writer thread WITHOUT flushing pending messages.
|
||||
|
||||
Used on shutdown when persistence is disabled (saveMessages: false):
|
||||
the thread must still be joined so process exit is clean, but nothing
|
||||
may be written.
|
||||
"""
|
||||
if self._async_queue is not None:
|
||||
if self._async_thread is not None and self._async_thread.is_alive():
|
||||
self._async_queue.put(_ASYNC_SHUTDOWN)
|
||||
self._async_thread.join(timeout=10)
|
||||
|
||||
def shutdown(self) -> None:
|
||||
"""Gracefully shut down the async writer thread."""
|
||||
if self._async_queue is not None:
|
||||
|
||||
@@ -316,6 +316,28 @@ class TestAsyncWriterThread:
|
||||
mgr.shutdown()
|
||||
assert mgr._async_thread is None
|
||||
|
||||
def test_stop_async_writer_joins_thread_without_flushing(self, make_manager):
|
||||
mgr = make_manager(write_frequency="async")
|
||||
mgr._ensure_async_writer()
|
||||
sess = _make_session()
|
||||
sess.add_message("user", "must not be written")
|
||||
with mgr._cache_lock:
|
||||
mgr._cache[sess.key] = sess
|
||||
|
||||
flushed = []
|
||||
mgr._flush_session = lambda session: flushed.append(session) or True
|
||||
|
||||
thread = mgr._async_thread
|
||||
mgr.stop_async_writer()
|
||||
thread.join(timeout=10)
|
||||
assert not thread.is_alive()
|
||||
assert flushed == []
|
||||
|
||||
def test_stop_async_writer_without_started_thread_is_noop(self, make_manager):
|
||||
mgr = make_manager(write_frequency="async")
|
||||
mgr.stop_async_writer()
|
||||
assert mgr._async_thread is None
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# async retry on failure
|
||||
|
||||
@@ -66,8 +66,10 @@ class TestOnSessionEnd:
|
||||
|
||||
|
||||
class TestShutdown:
|
||||
"""shutdown() joins worker threads then flushes; saveMessages=false must
|
||||
skip the flush (persistence) while still running the joins (cleanup)."""
|
||||
"""shutdown() joins worker threads then delegates to the session manager:
|
||||
manager.shutdown() (flush + join async writer) when persistence is on,
|
||||
manager.stop_async_writer() (join only, no flush) when saveMessages=false.
|
||||
Cleanup runs in both cases; only persistence is gated."""
|
||||
|
||||
def _provider_for_shutdown(self, save_messages: bool) -> HonchoMemoryProvider:
|
||||
p = _provider(save_messages=save_messages)
|
||||
@@ -78,12 +80,16 @@ class TestShutdown:
|
||||
p._sync_thread = None
|
||||
return p
|
||||
|
||||
def test_disabled_skips_flush(self):
|
||||
def test_disabled_skips_flush_but_stops_writer(self):
|
||||
p = self._provider_for_shutdown(save_messages=False)
|
||||
p.shutdown()
|
||||
p._manager.flush_all.assert_not_called()
|
||||
p._manager.shutdown.assert_not_called()
|
||||
p._manager.stop_async_writer.assert_called_once()
|
||||
|
||||
def test_enabled_flushes(self):
|
||||
def test_enabled_shuts_down_manager(self):
|
||||
p = self._provider_for_shutdown(save_messages=True)
|
||||
p.shutdown()
|
||||
p._manager.flush_all.assert_called_once()
|
||||
# manager.shutdown() flushes AND joins the async-writer thread;
|
||||
# calling flush_all() alone left the writer thread alive at exit.
|
||||
p._manager.shutdown.assert_called_once()
|
||||
|
||||
Reference in New Issue
Block a user