Skip to content

Commit 7f12e9c

Browse files
committed
Fix ProcessPoolExecutor shutdown deadlock when result pipe is full
1 parent 0906d2a commit 7f12e9c

2 files changed

Lines changed: 47 additions & 0 deletions

File tree

‎Lib/concurrent/futures/process.py‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -678,9 +678,17 @@ def _join_executor_internals(self, broken=False):
678678

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

686694
def get_n_children_alive(self):

‎Lib/test/test_concurrent_futures/test_deadlock.py‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
import queue
33
import signal
44
import sys
5+
import threading
56
import time
67
import unittest
78
import unittest.mock
@@ -329,6 +330,44 @@ def clear(self):
329330
signal.signal(signal.SIGALRM, old_handler)
330331

331332

333+
@warnings_helper.ignore_fork_in_thread_deprecation_warnings()
334+
def test_shutdown_should_not_deadlock_if_result_pipe_full(self):
335+
# Exiting workers put their pid on the result queue while shutdown()
336+
# joins them. If nothing drains the queue, a full pipe blocks the
337+
# workers forever. Fill the pipe once the executor manager thread
338+
# stops reading results, right before the workers are told to exit.
339+
self.executor.shutdown(wait=True)
340+
341+
data = b"a" * support.PIPE_MAX_SIZE
342+
fillers = []
343+
shutdown_workers = futures.process._ExecutorManagerThread.shutdown_workers
344+
def mock_shutdown_workers(self):
345+
# The put() blocks holding the result queue's write lock until
346+
# the pipe is drained, so every exiting worker blocks behind it.
347+
filler = threading.Thread(target=self.result_queue.put,
348+
args=(data,))
349+
filler.start()
350+
fillers.append(filler)
351+
shutdown_workers(self)
352+
353+
executor = self.executor_type(max_workers=2,
354+
mp_context=self.get_context())
355+
self.executor = executor # Allow clean up in fail_on_deadlock
356+
with unittest.mock.patch.object(futures.process._ExecutorManagerThread,
357+
'shutdown_workers',
358+
mock_shutdown_workers):
359+
self.assertEqual(list(executor.map(int, range(10))),
360+
list(range(10)))
361+
shutdown = threading.Thread(target=executor.shutdown)
362+
shutdown.start()
363+
shutdown.join(self.TIMEOUT)
364+
if shutdown.is_alive():
365+
self._fail_on_deadlock(executor)
366+
367+
self.assertEqual(len(fillers), 1)
368+
fillers[0].join()
369+
370+
332371
create_executor_tests(globals(), ExecutorDeadlockTest,
333372
executor_mixins=(ProcessPoolForkMixin,
334373
ProcessPoolForkserverMixin,

0 commit comments

Comments
 (0)