diff --git a/cron/scheduler.py b/cron/scheduler.py index bc3f5e7182..552fa848c7 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -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 diff --git a/tests/cron/test_sessiondb_init_hang.py b/tests/cron/test_sessiondb_init_hang.py index ce959325c5..924b3cba50 100644 --- a/tests/cron/test_sessiondb_init_hang.py +++ b/tests/cron/test_sessiondb_init_hang.py @@ -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)" + )