diff --git a/hermes_state.py b/hermes_state.py index 18c8e1aec7..0be368c618 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -21,6 +21,7 @@ import hashlib import json import logging import os +import queue import random import re import sqlite3 @@ -333,6 +334,11 @@ T = TypeVar("T") DEFAULT_DB_PATH = get_hermes_home() / "state.db" +# How long SessionDB stops attempting read-only opens after one fails, before +# probing again. Long enough that a genuinely unreadable file isn't retried per +# query; short enough that transient fd pressure doesn't strand the read pool. +_READ_OPEN_RETRY_SECONDS = 60.0 + # Import-time snapshot used by _default_db_path() to detect a deliberately # re-pointed DEFAULT_DB_PATH (tests monkeypatch the constant directly). _IMPORT_DEFAULT_DB_PATH = DEFAULT_DB_PATH @@ -2517,20 +2523,43 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) self.read_only = read_only self._lock = threading.Lock() - # Read-path split (WAL only): recall/browse queries run on per-thread - # read-only connections so they never queue behind writer flushes on - # self._lock. See _read_ctx(). - self._read_local = threading.local() - # Strong set of all live read connections across all threads. We - # hold a reference so short-lived reader threads' connections are - # not GC'd without close() — that would leak tracked fds in - # _live_connections. close() drains this set. - self._read_conns: "set[sqlite3.Connection]" = set() + # Read-path split (WAL only): recall/browse queries borrow a + # read-only connection from a bounded pool so they never queue + # behind writer flushes on self._lock. See _read_ctx(). + # + # The pool is BOUNDED because the previous per-thread + # (threading.local + strong set) scheme pinned one connection per + # (SessionDB x thread) for the life of the process. Starlette + # dispatches sync routes on anyio worker threads, so a SessionDB + # that is never closed accumulated a connection — and two fds, the + # database and its -wal — for every worker thread that ever read, + # until the process hit the 256 soft RLIMIT_NOFILE a service manager + # hands it and every request failed with EMFILE while the process + # stayed alive, so the supervisor's restart-on-exit never fired. + # Same bug class as the closing(...) fix in gateway/readiness.py + # (#69678 / #69567). + self._read_pool: "queue.LifoQueue[sqlite3.Connection]" = queue.LifoQueue( + maxsize=8 + ) self._read_conns_lock = threading.Lock() - # Set when close() begins. _get_read_conn checks this under the - # lock so a reader that finishes opening after the drain finds the - # shutdown in progress and closes its own connection immediately. + # Set when close() begins. _read_ctx checks this under the lock + # before returning a connection to the pool, so a reader still in + # flight during the drain closes its own connection instead of + # re-populating a pool nobody will drain again. self._read_conns_closed = False + # "read-only opens are failing against this file" backoff stamp. + # Instance-wide rather than per-thread: with a shared pool the open + # is no longer a per-thread event, and retrying a known-bad open on + # every query is a syscall storm for no benefit. The locked writer + # connection still serves reads while the backoff holds. + # Deliberately a TIMESTAMP, not a sticky bool: the likeliest trigger + # is transient fd pressure (EMFILE) — the very condition this pool + # exists to prevent — and a permanent flag would demote every reader + # on this instance to the writer lock for the life of the process. + # The gateway shares one SessionDB across every agent, so that turns + # a momentary blip into a permanent global convoy. Expires after + # _READ_OPEN_RETRY_SECONDS so the read path self-heals. + self._read_open_failed_at = 0.0 self._wal_active = False self._write_count = 0 # One-shot guard for the runtime FTS rebuild recovery on the write @@ -2762,7 +2791,10 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) # ── Read-path split ── def _get_read_conn(self) -> Optional[sqlite3.Connection]: - """Per-thread read-only connection, or None when unavailable. + """Open a fresh read-only connection, or None when unavailable. + + Callers must return the connection to self._read_pool (see + _read_ctx); this opens, it does not track. Only used under WAL: WAL readers see a consistent snapshot and never block on (or get blocked by) the writer, so recall/browse queries can @@ -2776,16 +2808,28 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) """ if not self._wal_active or self.read_only: return None - conn = getattr(self._read_local, "conn", None) - if conn is not None: - return conn - if getattr(self._read_local, "failed", False): - return None + with self._read_conns_lock: + if self._read_conns_closed: + return None + if ( + self._read_open_failed_at + and time.monotonic() - self._read_open_failed_at + < _READ_OPEN_RETRY_SECONDS + ): + return None try: conn = _connect_tracked_db( f"file:{self.db_path}?mode=ro", tracking_path=self.db_path, uri=True, + # Pooled connections are borrowed by whichever thread runs + # the next read, and sqlite3 otherwise refuses cross-thread + # use ("SQLite objects created in a thread can only be used + # in that same thread") — including on close(), which is how + # the old per-thread connections became unclosable and leaked + # their fds. Exclusive ownership is enforced by the pool + # checkout/return, not by sqlite3. Matches the writer opens. + check_same_thread=False, timeout=5.0, isolation_level=None, ) @@ -2797,36 +2841,75 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) # registry, not the database file, so mode=ro is fine. if self._fts_cjk_loaded: load_fts5_cjk_extension(conn) - with self._read_conns_lock: - if self._read_conns_closed: - # close() already drained — don't register; close - # immediately so no tracked fd leaks. - conn.close() - self._read_local.failed = True - return None - self._read_conns.add(conn) except sqlite3.Error: - # Mark this thread failed so we don't retry the open on every - # query; the locked writer connection still serves reads. - self._read_local.failed = True + # Back off from retrying the open on every query; the locked + # writer connection still serves reads until the stamp expires. + with self._read_conns_lock: + self._read_open_failed_at = time.monotonic() logger.debug("read-only connection open failed for %s", self.db_path, exc_info=True) return None - self._read_local.conn = conn return conn + def _close_read_conn(self, conn) -> None: + """Close a pooled read connection, reporting failures. + + This was a bare ``except Exception: pass``, which silently swallowed + the sqlite3.ProgrammingError raised when close() ran on a thread + other than the one that opened the connection — the exact signature + of the fd leak this pool fixes. A close that fails leaks a tracked + fd, so it must not be invisible. + """ + try: + conn.close() + except Exception as exc: + logger.warning("read-conn close failed for %s: %s", self.db_path, exc) + + def _checkout_read_conn(self) -> Optional[sqlite3.Connection]: + """Borrow a read connection from the pool, opening one on a miss. + + The single acquisition seam for the read path: the WAL/read_only gate, + the pool checkout and the open-on-miss all live here, so there is + exactly one place to exercise (and one place for a caller to bypass by + accident). Returns None when the read path is unavailable and the + caller must fall back to the locked writer connection. + """ + if not self._wal_active or self.read_only: + return None + try: + return self._read_pool.get_nowait() + except queue.Empty: + return self._get_read_conn() + @contextmanager def _read_ctx(self): """Yield a connection for read-only statements. - WAL: a per-thread read-only connection with NO lock — recall queries - never convoy behind writer flushes (the gateway shares one SessionDB - across every agent, so this lock was a global choke point). + WAL: a read-only connection borrowed from a bounded pool with NO + lock — recall queries never convoy behind writer flushes (the + gateway shares one SessionDB across every agent, so this lock was a + global choke point). The connection is checked out for the duration + of the block, so no two threads ever touch it concurrently. Non-WAL or read-conn failure: the shared writer connection under self._lock, byte-for-byte the legacy behavior. """ - conn = self._get_read_conn() + conn = self._checkout_read_conn() if conn is not None: - yield conn + try: + yield conn + finally: + returned = False + with self._read_conns_lock: + if not self._read_conns_closed: + try: + self._read_pool.put_nowait(conn) + returned = True + except queue.Full: + pass + if not returned: + # More concurrent readers than maxsize, or close() has + # already drained: this connection is surplus. Close it + # here — dropping it on the floor is what leaked the fd. + self._close_read_conn(conn) return with self._lock: yield self._conn @@ -3395,23 +3478,18 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin) # (instance, function), so this removes exactly our registration; # no-op when the writer never started. atexit.unregister(self._drain_token_queue_at_exit) - # Close all read-only connections across all threads. Per-thread - # connections live in threading.local() and would otherwise be GC'd - # without calling close(), leaking tracked fds in _live_connections. - # The strong set holds references so short-lived reader threads' - # connections survive until close() drains them. Setting the closed - # flag under the lock prevents a reader from registering a new - # connection after the drain. + # Drain the read-only connection pool. Setting the closed flag + # under the lock first means a reader still in flight closes its own + # connection on release instead of re-populating a pool that has + # already been drained. with self._read_conns_lock: self._read_conns_closed = True - read_conns = list(self._read_conns) - self._read_conns.clear() - for conn in read_conns: + while True: try: - conn.close() - except Exception: - pass - self._read_local.conn = None + conn = self._read_pool.get_nowait() + except queue.Empty: + break + self._close_read_conn(conn) with self._lock: if self._conn: if not self.read_only: diff --git a/tests/test_session_db_read_conn_pool.py b/tests/test_session_db_read_conn_pool.py new file mode 100644 index 0000000000..176c4d2419 --- /dev/null +++ b/tests/test_session_db_read_conn_pool.py @@ -0,0 +1,227 @@ +"""The SessionDB read path must not leak one connection per (SessionDB x thread). + +``_get_read_conn`` used to cache a read-only connection in ``threading.local()`` +and pin it in a strong set (``_read_conns``) that was only ever drained by +``close()``. Starlette dispatches sync routes on anyio worker threads, so a +SessionDB that is never closed -- the dashboard's module-global ``_db`` and +the per-session ``session_db`` handles -- gained a connection, and a file +descriptor, for every worker thread that ever served a read. In production +that walked into the 256 soft ``RLIMIT_NOFILE`` a service manager hands the +process, after which every request failed with ``OSError`` EMFILE while the +process stayed alive, so the supervisor's restart-on-exit never fired. + +Worse, those connections were opened WITHOUT ``check_same_thread=False`` (both +writer opens pass it), so ``close()`` on them raised ``ProgrammingError`` from +a different thread and the bare ``except Exception: pass`` hid it -- leaving +``hermes_cli.sqlite_safe_read``'s registry permanently over-counted as well. + +The contract pinned here: reads borrow from a BOUNDED pool, connections are +returned and reused, surplus connections are closed rather than dropped, and +``close()`` actually closes them from whatever thread it runs on. + +These assert on the pool/registry counts, never on ``lsof``: SQLite's unix VFS +parks a closed descriptor on a per-inode reuse list while any connection still +holds POSIX locks on that inode, so raw descriptor counts lag the real +connection count and make such assertions flaky. +""" + +import threading + +import pytest + +from hermes_state import SessionDB + + +def _live_count(path) -> int: + """Live-connection count the tracking registry holds for *path*.""" + import hermes_cli.sqlite_safe_read as mod + + with mod._live_lock: + return mod._live_connections.get(mod._key(path), 0) + + +@pytest.fixture() +def db(tmp_path): + d = SessionDB(db_path=tmp_path / "state.db") + d.create_session(session_id="s1", source="cli", model="m") + d.append_message("s1", role="user", content="hello graphiti world") + d.append_message("s1", role="assistant", content="the neo4j daemon is healthy") + yield d + d.close() + + +def _read(db): + db.get_session("s1") + db.search_messages("graphiti", limit=5) + db.get_messages("s1") + + +@pytest.mark.requires_wal +def test_read_pool_is_bounded_across_many_threads(db): + """150 short-lived reader threads must not pin 150 connections.""" + maxsize = db._read_pool.maxsize + assert maxsize > 0, "read pool must be bounded" + + for _ in range(6): + threads = [threading.Thread(target=_read, args=(db,)) for _ in range(25)] + for t in threads: + t.start() + for t in threads: + t.join() + assert db._read_pool.qsize() <= maxsize + + # The pre-fix code held 151 connections here (150 readers + main thread). + assert db._read_pool.qsize() <= maxsize + # +1 for the writer connection SessionDB always holds. + assert _live_count(db.db_path) <= maxsize + 1 + + +@pytest.mark.requires_wal +def test_read_conn_returned_to_pool_and_reused(db): + """Sequential reads on one thread reuse a pooled connection, not a new one.""" + with db._read_ctx() as conn: + first = conn + assert db._read_pool.qsize() >= 1, "connection was not returned to the pool" + with db._read_ctx() as conn: + assert conn is first, "pooled connection was not reused" + + +@pytest.mark.requires_wal +def test_pooled_conn_is_usable_from_another_thread(db): + """A pooled connection is handed between threads, so it must not be + bound to its creating thread (check_same_thread=False).""" + with db._read_ctx() as conn: + borrowed = conn + + errors = [] + + def use_it(): + try: + borrowed.execute("SELECT 1").fetchone() + except Exception as exc: # noqa: BLE001 + errors.append(exc) + + t = threading.Thread(target=use_it) + t.start() + t.join() + assert not errors, f"pooled connection unusable off-thread: {errors}" + + +@pytest.mark.requires_wal +def test_close_drains_pool_from_a_foreign_thread(tmp_path): + """close() must actually close pooled connections, including ones opened + on threads that have since exited -- the swallowed ProgrammingError.""" + d = SessionDB(db_path=tmp_path / "state2.db") + d.create_session(session_id="s1", source="cli", model="m") + + # Populate the pool from a worker thread, then let that thread die. + t = threading.Thread(target=lambda: d.get_session("s1")) + t.start() + t.join() + assert d._read_pool.qsize() >= 1 + + d.close() + assert d._read_pool.qsize() == 0 + # Registry back to zero proves the closes succeeded rather than raising + # ProgrammingError into a bare except. + assert _live_count(d.db_path) == 0 + + +@pytest.mark.requires_wal +def test_reader_after_close_does_not_repopulate_pool(db): + """A read racing close() must close its connection, not refill the pool.""" + db.close() + assert db._read_pool.qsize() == 0 + # A read arriving after the drain must not open-and-requeue a connection + # that nothing will ever close again. + with db._read_ctx(): + pass + assert db._read_pool.qsize() == 0 + + +def test_reads_are_still_correct_under_concurrency(db): + """Pooling must not corrupt results when threads share connections.""" + results = [] + errors = [] + + def reader(): + try: + results.append(db.get_session("s1")["id"]) + results.append(len(db.get_messages("s1"))) + except Exception as exc: # noqa: BLE001 + errors.append(exc) + + threads = [threading.Thread(target=reader) for _ in range(12)] + for t in threads: + t.start() + for t in threads: + t.join() + assert not errors, f"concurrent reads failed: {errors}" + assert results.count("s1") == 12 + assert results.count(2) == 12 + + +@pytest.mark.requires_wal +def test_read_open_failure_backs_off_but_recovers(db): + """A failed read-only open must not permanently demote the read path. + + The first version of this fix used a sticky instance-wide boolean + (``_read_open_failed``). Its likeliest trigger is transient fd pressure -- + EMFILE, the very condition this pool exists to prevent -- and because the + gateway shares ONE SessionDB across every agent, a single blip would have + convoyed every subsequent reader behind the writer lock for the life of + the process. The stamp must expire. + """ + import time as _time + + from hermes_state import _READ_OPEN_RETRY_SECONDS + + baseline = db._get_read_conn() + assert baseline is not None, "baseline read open should succeed" + db._close_read_conn(baseline) + + db._read_open_failed_at = _time.monotonic() + assert db._get_read_conn() is None, "should back off immediately after a failure" + + db._read_open_failed_at = _time.monotonic() - (_READ_OPEN_RETRY_SECONDS + 1) + recovered = db._get_read_conn() + assert recovered is not None, "read path must self-heal once the window expires" + db._close_read_conn(recovered) + + +@pytest.mark.requires_wal +def test_checkout_seam_is_the_single_acquisition_point(db): + """``_read_ctx`` must acquire via ``_checkout_read_conn`` and nothing else. + + If a future edit re-inlines the pool checkout into ``_read_ctx``, patching + ``_get_read_conn`` silently exercises nothing whenever the pool is warm -- + which is exactly how the writer-lock fallback test below would rot into a + no-op without failing. + """ + calls = [] + original = db._checkout_read_conn + + def _spy(): + calls.append(1) + return original() + + db._checkout_read_conn = _spy + try: + with db._read_ctx(): + pass + finally: + db._checkout_read_conn = original + assert calls, "_read_ctx must route acquisition through _checkout_read_conn" + + +def test_fallback_to_locked_writer_when_read_conn_unavailable(db, monkeypatch): + """With no read connection available, reads still work under self._lock. + + Patched at the acquisition SEAM rather than at ``_get_read_conn``: the + pool is consulted first, so a patched ``_get_read_conn`` is never reached + while the pool holds a connection and this test would pass while + exercising nothing. + """ + monkeypatch.setattr(db, "_checkout_read_conn", lambda: None) + assert db.get_session("s1")["id"] == "s1" + assert db.search_messages("graphiti", limit=5) diff --git a/tests/test_session_db_read_path_split.py b/tests/test_session_db_read_path_split.py index 1b28b90589..1f3095ca87 100644 --- a/tests/test_session_db_read_path_split.py +++ b/tests/test_session_db_read_path_split.py @@ -1,11 +1,12 @@ -"""Tests for the SessionDB read-path split (per-thread read-only connections). +"""Tests for the SessionDB read-path split (pooled read-only connections). The gateway shares ONE SessionDB across every agent, so recall/browse reads used to queue behind writer flushes on self._lock — a measured production convoy (a 0.2s FTS query stretched to 112s while 6-8 concurrent turns flushed tool results). These tests pin the new contract: reads run on a -per-thread read-only connection under WAL, never touch self._lock, and fall -back to the legacy locked path when WAL or the read connection is missing. +read-only connection borrowed from a bounded pool under WAL, never touch +self._lock, and fall back to the legacy locked path when WAL or the read +connection is missing. """ import threading @@ -39,8 +40,19 @@ def test_read_conn_is_per_thread(db): assert conns[1] is not conns[2] -def test_read_conn_reused_within_thread(db): - assert db._get_read_conn() is db._get_read_conn() +@pytest.mark.requires_wal +def test_read_conn_reused_via_pool(db): + """Reuse is now the pool's job, not a per-thread memo. + + The old contract (``_get_read_conn()`` returns the same object twice on one + thread) was the leak: that memo pinned one unclosable connection per + (SessionDB x thread) forever. ``_get_read_conn`` now always opens a fresh + connection and reuse happens via checkout/return, so assert on that. + """ + with db._read_ctx() as first: + assert first is not None + with db._read_ctx() as second: + assert second is first, "sequential readers must reuse the pooled conn" @pytest.mark.requires_wal