fix(cron): close leaked SessionDB connection when init outlives the timeout-abandoned worker (#72782)
run_job() submits SessionDB() to a one-worker executor and abandons the worker (shutdown(wait=False)) when init exceeds the cron timeout. If the constructor later completes inside that abandoned worker, the Future's result — an open SessionDB holding .db/WAL/SHM handles — was orphaned and never closed, leaking descriptors until EMFILE. Attach a done-callback on the timeout path that retrieves and closes any eventual late result. Salvage note: the lazy-recall ownership half of #72822 (_owns_session_db tracked on AIAgent, owned handle closed in close()) already landed on main; this carries the remaining cron timeout-abandon half with its regression test.
This commit is contained in:
+28
-1
@@ -63,6 +63,24 @@ from agent.delegation_context import (
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _close_late_session_db_result(future: "concurrent.futures.Future") -> None:
|
||||
"""Done-callback: close a SessionDB whose constructor finished after run_job's timeout.
|
||||
|
||||
When ``run_job``'s SessionDB init times out, the worker thread is abandoned
|
||||
(``shutdown(wait=False)``) so the job can proceed without a session store.
|
||||
If the constructor later completes inside that abandoned worker, the
|
||||
Future's result — an open SessionDB holding .db / WAL / SHM file handles —
|
||||
would be orphaned and never closed, leaking descriptors until EMFILE
|
||||
(#72782). This callback retrieves and closes that eventual late result.
|
||||
"""
|
||||
try:
|
||||
db = future.result()
|
||||
if db is not None:
|
||||
db.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _set_cron_session_title(session_db, session_id, base_title):
|
||||
"""Robustly title a finished cron session before it is closed.
|
||||
|
||||
@@ -4186,8 +4204,17 @@ def run_job(
|
||||
|
||||
if _session_db_timeout > 0:
|
||||
_session_db_pool = concurrent.futures.ThreadPoolExecutor(max_workers=1)
|
||||
_session_db_future = _session_db_pool.submit(SessionDB)
|
||||
try:
|
||||
_session_db = _session_db_pool.submit(SessionDB).result(timeout=_session_db_timeout)
|
||||
_session_db = _session_db_future.result(timeout=_session_db_timeout)
|
||||
except concurrent.futures.TimeoutError:
|
||||
# The worker is abandoned (shutdown below doesn't wait for it).
|
||||
# If SessionDB() later completes inside it, the future's result
|
||||
# would be orphaned and its SQLite FDs (.db, WAL, SHM) leak
|
||||
# until process exit. Register a done-callback that retrieves
|
||||
# and closes any eventual late result (#72782).
|
||||
_session_db_future.add_done_callback(_close_late_session_db_result)
|
||||
raise
|
||||
finally:
|
||||
# Don't wait for a wedged connect() to unwind — abandon the
|
||||
# worker thread (same pattern as the agent inactivity timeout
|
||||
|
||||
@@ -245,3 +245,101 @@ class TestDispatchGuardReleasedAfterHang:
|
||||
finally:
|
||||
sched._running_job_ids.discard("guard-sessiondb-hang")
|
||||
sched._shutdown_parallel_pool()
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# Bug #72782: late SessionDB result leaks FDs after timeout abandonment
|
||||
# ===========================================================================
|
||||
|
||||
class TestCloseLateSessionDbResult:
|
||||
"""Unit tests for the done-callback that closes a SessionDB whose
|
||||
constructor completed after run_job's timeout."""
|
||||
|
||||
def test_closes_db_from_completed_future(self):
|
||||
"""A completed future holding a SessionDB is closed."""
|
||||
import concurrent.futures
|
||||
from cron.scheduler import _close_late_session_db_result
|
||||
|
||||
mock_db = MagicMock()
|
||||
fut = concurrent.futures.Future()
|
||||
fut.set_result(mock_db)
|
||||
|
||||
_close_late_session_db_result(fut)
|
||||
|
||||
mock_db.close.assert_called_once()
|
||||
|
||||
def test_safe_when_result_is_none(self):
|
||||
"""No error when the future's result is None."""
|
||||
import concurrent.futures
|
||||
from cron.scheduler import _close_late_session_db_result
|
||||
|
||||
fut = concurrent.futures.Future()
|
||||
fut.set_result(None)
|
||||
_close_late_session_db_result(fut) # must not raise
|
||||
|
||||
def test_safe_when_future_raised(self):
|
||||
"""No error when the future itself raised (e.g. connect failed)."""
|
||||
import concurrent.futures
|
||||
from cron.scheduler import _close_late_session_db_result
|
||||
|
||||
fut = concurrent.futures.Future()
|
||||
fut.set_exception(RuntimeError("connect failed"))
|
||||
_close_late_session_db_result(fut) # must not raise
|
||||
|
||||
|
||||
class TestLateSessionDbClosedAfterTimeout:
|
||||
"""End-to-end: when SessionDB init times out but later completes inside the
|
||||
abandoned worker, the orphaned result must be closed (#72782)."""
|
||||
|
||||
def test_late_session_db_result_is_closed(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("HERMES_CRON_SESSION_DB_TIMEOUT", "0.2")
|
||||
never_set = threading.Event()
|
||||
late_db_holder = [] # captures the SessionDB returned by the late init
|
||||
|
||||
def _hanging_then_capture():
|
||||
never_set.wait(timeout=30)
|
||||
db = MagicMock()
|
||||
late_db_holder.append(db)
|
||||
return db
|
||||
|
||||
job = {"id": "late-close-test", "name": "test", "prompt": "hello"}
|
||||
|
||||
try:
|
||||
with patch("cron.scheduler._hermes_home", tmp_path), \
|
||||
patch("cron.scheduler._resolve_origin", return_value=None), \
|
||||
patch("hermes_cli.env_loader.load_hermes_dotenv"), \
|
||||
patch("hermes_cli.env_loader.reset_secret_source_cache"), \
|
||||
patch("hermes_state.SessionDB", side_effect=_hanging_then_capture), \
|
||||
patch(
|
||||
"hermes_cli.runtime_provider.resolve_runtime_provider",
|
||||
return_value={
|
||||
"api_key": "test-key",
|
||||
"base_url": "https://example.invalid/v1",
|
||||
"provider": "openrouter",
|
||||
"api_mode": "chat_completions",
|
||||
},
|
||||
), \
|
||||
patch("run_agent.AIAgent") as mock_agent_cls:
|
||||
mock_agent = MagicMock()
|
||||
mock_agent.run_conversation.return_value = {"final_response": "ok"}
|
||||
mock_agent_cls.return_value = mock_agent
|
||||
|
||||
success, output, final_response, error = run_job(job)
|
||||
# run_job returned promptly after the timeout; session_db is None
|
||||
assert success is True
|
||||
|
||||
# Release the hanging init so the abandoned worker completes.
|
||||
never_set.set()
|
||||
# Wait for the done-callback to fire and close the late result.
|
||||
for _ in range(50):
|
||||
if late_db_holder and late_db_holder[0].close.called:
|
||||
break
|
||||
time.sleep(0.1)
|
||||
finally:
|
||||
never_set.set()
|
||||
|
||||
assert len(late_db_holder) == 1, "SessionDB() should have completed once"
|
||||
late_db_holder[0].close.assert_called_once(), (
|
||||
"The SessionDB that completed after the timeout must be closed by "
|
||||
"the done-callback — otherwise its SQLite FDs leak until process exit (#72782)"
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user