fix: a failed executor submission releases the API worker count (review follow-up)
The count was incremented on the submitting thread and released only in the worker's finally; when loop.run_in_executor itself raised (default executor already shut down during quiesce -> RuntimeError) no worker ever ran and _API_WORKER_LIVE stayed elevated for the process lifetime, making the shutdown close gate skip the SessionDB close forever. _submit_api_worker now owns both sides of the submission.
This commit is contained in:
@@ -4048,7 +4048,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
try:
|
||||
# Worker-scoped count rides along so the shutdown close gate still sees the thread
|
||||
# after this handler task is cancelled (#116535); released in the worker's finally.
|
||||
return await loop.run_in_executor(None, _api_runs._track_api_worker(_run))
|
||||
return await _api_runs._submit_api_worker(loop, _run)
|
||||
finally:
|
||||
self._inflight_agent_runs -= 1
|
||||
|
||||
|
||||
@@ -45,12 +45,14 @@ def api_worker_live_count() -> int:
|
||||
return _API_WORKER_LIVE
|
||||
|
||||
|
||||
def _track_api_worker(fn):
|
||||
"""Hold the worker-lifetime count for one ``run_in_executor`` submission.
|
||||
def _submit_api_worker(loop, fn):
|
||||
"""``loop.run_in_executor(None, fn)`` with the worker-lifetime count held for the submission.
|
||||
|
||||
Increment on the submitting (handler) thread so the count is live before the worker can
|
||||
exit; decrement in the worker thread's own ``finally`` so handler cancellation cannot drop
|
||||
it early. ``fn`` still runs entirely on the worker.
|
||||
it early. When submission itself fails (default executor already shut down during quiesce
|
||||
-> RuntimeError) the worker never runs, so the count is released here instead — a leaked
|
||||
count would make the shutdown close gate skip the SessionDB close for the process lifetime.
|
||||
"""
|
||||
global _API_WORKER_LIVE
|
||||
with _API_WORKER_LOCK:
|
||||
@@ -64,7 +66,12 @@ def _track_api_worker(fn):
|
||||
with _API_WORKER_LOCK:
|
||||
_API_WORKER_LIVE -= 1
|
||||
|
||||
return _counted
|
||||
try:
|
||||
return loop.run_in_executor(None, _counted)
|
||||
except BaseException:
|
||||
with _API_WORKER_LOCK:
|
||||
_API_WORKER_LIVE -= 1
|
||||
raise
|
||||
_ROOM_RETENTION_REQUEST_KEY = (
|
||||
RequestKey("hermes.room_run_retention_until", float) if RequestKey is not None
|
||||
else "hermes.room_run_retention_until")
|
||||
@@ -905,9 +912,8 @@ async def _execute_run(self, run: _RunLaunch, *, _api_server) -> None:
|
||||
interim_assistant_callback=_interim_cb, **run.agent_kwargs)
|
||||
self._active_run_agents[run_id] = agent
|
||||
approval_notify = _make_approval_notify(self, run, _api_server=_api_server)
|
||||
result, usage, served_runtime = await loop.run_in_executor(
|
||||
None, _track_api_worker(
|
||||
lambda: _run_agent_sync(self, run, agent, approval_notify, _api_server=_api_server)))
|
||||
result, usage, served_runtime = await _submit_api_worker(
|
||||
loop, lambda: _run_agent_sync(self, run, agent, approval_notify, _api_server=_api_server))
|
||||
if not isinstance(result, dict):
|
||||
result = {}
|
||||
status, fields = terminal_run_status(result)
|
||||
|
||||
@@ -726,3 +726,18 @@ class TestShutdownSettleWindow:
|
||||
_INTERRUPT_REASON_GATEWAY_SHUTDOWN,
|
||||
]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_failed_executor_submission_releases_the_worker_count():
|
||||
"""A request that reaches ``run_in_executor`` after ``shutdown_default_executor()`` raises
|
||||
RuntimeError and never runs a worker; the worker-scoped count must not stay elevated for the
|
||||
process lifetime, or the shutdown SessionDB-close gate skips the close forever (#116535)."""
|
||||
baseline = _api_runs.api_worker_live_count()
|
||||
|
||||
class _ShutExecutorLoop:
|
||||
def run_in_executor(self, executor, fn):
|
||||
raise RuntimeError("Executor shutdown has been called")
|
||||
|
||||
with pytest.raises(RuntimeError, match="Executor shutdown"):
|
||||
_api_runs._submit_api_worker(_ShutExecutorLoop(), lambda: None)
|
||||
assert _api_runs.api_worker_live_count() == baseline
|
||||
|
||||
@@ -303,7 +303,7 @@ async def test_cancelled_api_handler_worker_still_blocks_session_db_close(monkey
|
||||
# Same shape as the api_server call sites: the worker-scoped count is taken before
|
||||
# run_in_executor and released in the worker's own finally, so cancelling this task
|
||||
# drops the handler side while the thread keeps holding the worker side.
|
||||
return await loop.run_in_executor(None, api_runs._track_api_worker(_blocked_turn))
|
||||
return await api_runs._submit_api_worker(loop, _blocked_turn)
|
||||
|
||||
task = asyncio.ensure_future(_handler())
|
||||
assert await loop.run_in_executor(None, worker_started.wait, 5.0), "worker never started"
|
||||
|
||||
Reference in New Issue
Block a user