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
8 changes: 8 additions & 0 deletions Lib/concurrent/futures/process.py
Original file line number Diff line number Diff line change
Expand Up @@ -678,9 +678,17 @@ def _join_executor_internals(self, broken=False):

# If .join() is not called on the created processes then
# some ctx.Queue methods may deadlock on Mac OS X.
result_reader = self.result_queue._reader
for p in self.processes.values():
if broken:
p.terminate()
else:
# Exiting workers put their pid on the result queue. Keep
# draining it, otherwise once the pipe is full a worker
# blocks forever writing to it and never exits.
while p.sentinel not in mp.connection.wait(
[result_reader, p.sentinel]):
result_reader.recv_bytes()
p.join()

def get_n_children_alive(self):
Expand Down
39 changes: 39 additions & 0 deletions Lib/test/test_concurrent_futures/test_deadlock.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
import queue
import signal
import sys
import threading
import time
import unittest
import unittest.mock
Expand Down Expand Up @@ -329,6 +330,44 @@ def clear(self):
signal.signal(signal.SIGALRM, old_handler)


@warnings_helper.ignore_fork_in_thread_deprecation_warnings()
def test_shutdown_should_not_deadlock_if_result_pipe_full(self):
# Exiting workers put their pid on the result queue while shutdown()
# joins them. If nothing drains the queue, a full pipe blocks the
# workers forever. Fill the pipe once the executor manager thread
# stops reading results, right before the workers are told to exit.
self.executor.shutdown(wait=True)

data = b"a" * support.PIPE_MAX_SIZE
fillers = []
shutdown_workers = futures.process._ExecutorManagerThread.shutdown_workers
def mock_shutdown_workers(self):
# The put() blocks holding the result queue's write lock until
# the pipe is drained, so every exiting worker blocks behind it.
filler = threading.Thread(target=self.result_queue.put,
args=(data,))
filler.start()
fillers.append(filler)
shutdown_workers(self)

executor = self.executor_type(max_workers=2,
mp_context=self.get_context())
self.executor = executor # Allow clean up in fail_on_deadlock
with unittest.mock.patch.object(futures.process._ExecutorManagerThread,
'shutdown_workers',
mock_shutdown_workers):
self.assertEqual(list(executor.map(int, range(10))),
list(range(10)))
shutdown = threading.Thread(target=executor.shutdown)
shutdown.start()
shutdown.join(self.TIMEOUT)
if shutdown.is_alive():
self._fail_on_deadlock(executor)

self.assertEqual(len(fillers), 1)
fillers[0].join()


create_executor_tests(globals(), ExecutorDeadlockTest,
executor_mixins=(ProcessPoolForkMixin,
ProcessPoolForkserverMixin,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
Fix a deadlock in :meth:`ProcessPoolExecutor.shutdown()
<concurrent.futures.Executor.shutdown>` when the messages that exiting
workers send back fill up the pipe used to return results. This could
happen with many workers or a small pipe buffer.
Loading