diff --git a/Lib/concurrent/futures/process.py b/Lib/concurrent/futures/process.py index 7f4f225c0ad4fbd..904787c081ee014 100644 --- a/Lib/concurrent/futures/process.py +++ b/Lib/concurrent/futures/process.py @@ -559,14 +559,18 @@ def _terminate_broken(self, cause, bpe_message=None): # lead the executor into becoming broken. bpe_message overrides the # default message on the BrokenProcessPool set on pending futures. - # Mark the process pool broken so that submits fail right now. - executor = self.executor_reference() - if executor is not None: - executor._broken = ('A child process terminated ' - 'abruptly, the process pool is not ' - 'usable anymore') - executor._shutdown_thread = True - executor = None + # Mark the pool broken while holding the lock so that concurrent + # calls to submit() fail instead of adding more work. Do not hold the + # lock while waiting for workers to exit below: a worker may ignore + # terminate(), and submit() and shutdown(wait=False) must remain + # responsive in that case. + with self.shutdown_lock: + executor = self.executor_reference() + if executor is not None: + executor._broken = ('A child process terminated ' + 'abruptly, the process pool is not ' + 'usable anymore') + executor._shutdown_thread = True # All pending tasks are to be marked failed with a # BrokenProcessPool error, as separate instances to avoid sharing @@ -619,8 +623,7 @@ def _terminate_broken(self, cause, bpe_message=None): self._join_executor_internals(broken=True) def terminate_broken(self, cause, bpe_message=None): - with self.shutdown_lock: - self._terminate_broken(cause, bpe_message) + self._terminate_broken(cause, bpe_message) def flag_executor_shutting_down(self): # Flag the executor as shutting down and cancel remaining tasks if diff --git a/Lib/test/test_concurrent_futures/test_process_pool.py b/Lib/test/test_concurrent_futures/test_process_pool.py index dafbda862c51c24..b0b1b34cb01324a 100644 --- a/Lib/test/test_concurrent_futures/test_process_pool.py +++ b/Lib/test/test_concurrent_futures/test_process_pool.py @@ -543,6 +543,64 @@ def test_force_shutdown_workers_stops_pool(self, function_name): break +class BrokenPoolCleanupTest(unittest.TestCase): + + def test_shutdown_lock_released_before_joining_workers(self): + join_started = threading.Event() + release_join = threading.Event() + + class Executor: + pass + + class Process: + pid = 1 + exitcode = None + + def terminate(self): + pass + + def join(self): + join_started.set() + release_join.wait() + + class CallQueue: + def _terminate_broken(self): + pass + + def close(self): + pass + + def join_thread(self): + pass + + class ThreadWakeup: + def close(self): + pass + + executor = Executor() + executor._broken = None + executor._shutdown_thread = False + manager = object.__new__(futures.process._ExecutorManagerThread) + manager.shutdown_lock = threading.Lock() + manager.executor_reference = weakref.ref(executor) + manager.processes = {1: Process()} + manager.pending_work_items = {} + manager.call_queue = CallQueue() + manager.thread_wakeup = ThreadWakeup() + manager_thread = threading.Thread( + target=manager.terminate_broken, args=(None,)) + manager_thread.start() + try: + self.assertTrue(join_started.wait(support.SHORT_TIMEOUT)) + self.assertTrue( + manager.shutdown_lock.acquire(timeout=support.SHORT_TIMEOUT)) + manager.shutdown_lock.release() + finally: + release_join.set() + manager_thread.join(support.SHORT_TIMEOUT) + self.assertFalse(manager_thread.is_alive()) + + create_executor_tests(globals(), ProcessPoolExecutorTest, executor_mixins=(ProcessPoolForkMixin, ProcessPoolForkserverMixin, diff --git a/Misc/NEWS.d/next/Library/2026-09-29-23-20-00.gh-issue-158413.aB93Lm.rst b/Misc/NEWS.d/next/Library/2026-09-29-23-20-00.gh-issue-158413.aB93Lm.rst new file mode 100644 index 000000000000000..b59ef8110d0e4f1 --- /dev/null +++ b/Misc/NEWS.d/next/Library/2026-09-29-23-20-00.gh-issue-158413.aB93Lm.rst @@ -0,0 +1,2 @@ +:class:`concurrent.futures.ProcessPoolExecutor` no longer holds its shutdown +lock while waiting for workers to exit after a worker fails.