Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 13 additions & 10 deletions Lib/concurrent/futures/process.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
58 changes: 58 additions & 0 deletions Lib/test/test_concurrent_futures/test_process_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
:class:`concurrent.futures.ProcessPoolExecutor` no longer holds its shutdown
lock while waiting for workers to exit after a worker fails.
Loading