state: pool SessionDB read connections instead of leaking one per (SessionDB x thread)

This commit is contained in:
Yishova
2026-08-02 05:23:34 -04:00
committed by Teknium
parent 6a7cf19302
commit 87aedbe7b6
3 changed files with 371 additions and 54 deletions

View File

@@ -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:

View File

@@ -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)

View File

@@ -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