fix(state): exclude closed SessionDB writers
This commit is contained in:
@@ -562,7 +562,7 @@ class SessionDB(
|
||||
# per DATABASE PATH, not per instance: the descriptors they ration belong to the file, and one
|
||||
# process holds several SessionDB objects on the same state.db (#98573). See _PathReadBudget.
|
||||
self._read_budget = _read_budget_for(self.db_path)
|
||||
self._read_budget.register(self)
|
||||
self._read_budget_registered = False
|
||||
self._read_permits = self._read_budget.permits
|
||||
self._read_conns_lock = threading.Lock()
|
||||
# Set when close() begins; an in-flight reader then closes its own connection
|
||||
@@ -623,6 +623,10 @@ class SessionDB(
|
||||
conn, self._conn = self._conn, None
|
||||
self._close_connection_quietly(conn)
|
||||
else:
|
||||
# Only a successfully opened handle owns a writer connection. Failed
|
||||
# construction must not leave a diagnostic member behind.
|
||||
self._read_budget.register(self)
|
||||
self._read_budget_registered = True
|
||||
# Test-isolation runs only (gated inside the helper): register
|
||||
# for the suite-level leak sweep in tests/conftest.py.
|
||||
_register_test_instance(self)
|
||||
@@ -1488,6 +1492,9 @@ class SessionDB(
|
||||
# Only a clean close ends the generation; retain the recorded
|
||||
# identity when retiring an unsafe handle.
|
||||
self._db_sidecar_identity = {}
|
||||
if self._read_budget_registered:
|
||||
self._read_budget.unregister(self)
|
||||
self._read_budget_registered = False
|
||||
|
||||
def __del__(self) -> None:
|
||||
"""Safety net: close() if the caller forgot. Attribute access stays
|
||||
|
||||
@@ -163,7 +163,10 @@ class _PathReadBudget:
|
||||
# Only writable handles carry the cost the warning names (writer connection, write
|
||||
# lock, close-time checkpoint). Read-only attaches (dashboard routers, status/lookup
|
||||
# one-shots) open per request by design and must not trip it.
|
||||
handles = sum(1 for member in self._members if not member.read_only)
|
||||
handles = sum(
|
||||
1 for member in self._members
|
||||
if not member.read_only and member._conn is not None
|
||||
)
|
||||
warn = (handles > _HANDLES_PER_PATH_WARN and not self._duplicate_handles_warned)
|
||||
if warn:
|
||||
self._duplicate_handles_warned = True
|
||||
@@ -187,6 +190,11 @@ class _PathReadBudget:
|
||||
handles, db.db_path, _READ_POOL_MAX, creation_sites,
|
||||
)
|
||||
|
||||
def unregister(self, db: "SessionDB") -> None:
|
||||
"""Remove a closed writer from duplicate-handle diagnostics immediately."""
|
||||
with self._lock:
|
||||
self._members.discard(db)
|
||||
|
||||
def acquire(self, requester: "SessionDB") -> bool:
|
||||
"""Take a permit for a new read connection, or refuse (caller degrades to the
|
||||
locked writer connection). Gates, broadest first: fd headroom, process
|
||||
|
||||
@@ -37,6 +37,7 @@ connection count and make such assertions flaky.
|
||||
|
||||
import hermes_state_readpool
|
||||
import queue
|
||||
import sqlite3
|
||||
import threading
|
||||
|
||||
import pytest
|
||||
@@ -740,3 +741,47 @@ def test_handle_diagnostics_unavailable_does_not_block_database(tmp_path, monkey
|
||||
handle.create_session(session_id="available", source="cli", model="m")
|
||||
assert handle.get_session("available")["id"] == "available"
|
||||
assert handle._creation_site == "unknown"
|
||||
|
||||
|
||||
@pytest.mark.requires_wal
|
||||
def test_closed_handles_do_not_count_toward_duplicate_writer_warning(db, caplog):
|
||||
"""Closing a writer releases its duplicate-writer diagnostic membership."""
|
||||
import logging
|
||||
|
||||
from hermes_state_readpool import _HANDLES_PER_PATH_WARN
|
||||
|
||||
closed = [SessionDB(db_path=db.db_path) for _ in range(_HANDLES_PER_PATH_WARN - 1)]
|
||||
for handle in closed:
|
||||
handle.close()
|
||||
|
||||
caplog.clear()
|
||||
with caplog.at_level(logging.WARNING, logger="hermes_state"):
|
||||
survivor = SessionDB(db_path=db.db_path)
|
||||
try:
|
||||
assert not any(
|
||||
"live SessionDB handles on" in record.getMessage() for record in caplog.records
|
||||
), "closed writers were retained as live duplicate handles"
|
||||
finally:
|
||||
survivor.close()
|
||||
|
||||
|
||||
def test_failed_initialization_does_not_register_duplicate_writer_handle(db, monkeypatch):
|
||||
"""A constructor that raises before opening must never join the handle budget."""
|
||||
budget = hermes_state_readpool._read_budget_for(db.db_path)
|
||||
registered = []
|
||||
original_register = budget.register
|
||||
|
||||
def remember_register(handle):
|
||||
registered.append(handle)
|
||||
original_register(handle)
|
||||
|
||||
def fail_open_writer(self):
|
||||
raise sqlite3.OperationalError("injected initialization failure")
|
||||
|
||||
monkeypatch.setattr(budget, "register", remember_register)
|
||||
monkeypatch.setattr(SessionDB, "_open_writer", fail_open_writer)
|
||||
|
||||
with pytest.raises(sqlite3.OperationalError, match="injected initialization failure"):
|
||||
SessionDB(db_path=db.db_path)
|
||||
|
||||
assert registered == []
|
||||
|
||||
Reference in New Issue
Block a user