gh-155477: Fix multiprocessing.Pool deadlock on close() with a buffersize imap (GH-155478)

serhiy-storchaka <[email protected]>
Newsgroups gmane.comp.python.cvs
Message-ID <[email protected]>
https://github.com/python/cpython/commit/70fdc966d0420e64e848eb09f7f695735f2a8b66
commit: 70fdc966d0420e64e848eb09f7f695735f2a8b66
branch: main
author: tonghuaroot (童话) <[email protected]>
committer: serhiy-storchaka <[email protected]>
date: 2026-08-16T13:43:24Z
summary:

gh-155477: Fix multiprocessing.Pool deadlock on close() with a buffersize imap (GH-155478)

close() did not release the buffersize semaphores that throttle the task
generator, so a partially-consumed imap left the task handler blocked and
join() deadlocked. Release them in close() and stop the generator once the
pool leaves the RUN state.

files:
M Lib/multiprocessing/pool.py
M Lib/test/_test_multiprocessing.py

diff --git a/Lib/multiprocessing/pool.py b/Lib/multiprocessing/pool.py
index 8fd0f98a02dd3a6..f50bcbe4451bea4 100644
--- a/Lib/multiprocessing/pool.py
+++ b/Lib/multiprocessing/pool.py
@@ -403,6 +403,11 @@ def _guarded_task_generation(self, result_job, func, iterable, sema=None):
                 enumerated_iter = iter(enumerate(iterable))
                 while True:
                     sema.acquire()
+                    if self._state != RUN:
+                        # The pool is closing or terminating; stop submitting
+                        # the still-throttled tasks so the task handler can
+                        # finish instead of blocking here forever.
+                        break
                     try:
                         i, x = next(enumerated_iter)
                     except StopIteration:
@@ -661,6 +666,10 @@ def close(self):
             self._state = CLOSE
             self._worker_handler._state = CLOSE
             self._change_notifier.put(None)
+            # Wake any task generator throttled on a buffersize semaphore so
+            # it observes the CLOSE state and stops submitting.
+            for sema in list(self._taskqueue_buffersize_semaphores):
+                sema.release()
 
     def terminate(self):
         util.debug('terminating pool')
diff --git a/Lib/test/_test_multiprocessing.py b/Lib/test/_test_multiprocessing.py
index 338a31fd7f869ea..e5f618f5f2e84f4 100644
--- a/Lib/test/_test_multiprocessing.py
+++ b/Lib/test/_test_multiprocessing.py
@@ -3229,6 +3229,27 @@ def produce_args():
         p.terminate()
         p.join()
 
+    @warnings_helper.ignore_fork_in_thread_deprecation_warnings()
+    @support.subTests('method_name', ("imap", "imap_unordered"))
+    def test_imap_with_buffersize_close_after_partial_consumption(
+        self, method_name
+    ):
+        # close()/join() must not deadlock when a buffersize iterator is
+        # only partially consumed (the throttled task generator must stop).
+        p = self.Pool(2)
+        method = getattr(p, method_name)
+        it = method(sqr, range(1000), buffersize=2)
+        next(it)
+        finished = threading.Event()
+        def finalize():
+            p.close()
+            p.join()
+            finished.set()
+        t = threading.Thread(target=finalize)
+        t.start()
+        t.join(support.SHORT_TIMEOUT)
+        self.assertTrue(finished.is_set(), "close()/join() deadlocked")
+
     @support.subTests('method_name', ("imap", "imap_unordered"))
     def test_imap_and_imap_unordered_with_buffersize_on_empty_iterable(
         self, method_name

_______________________________________________
Python-checkins mailing list -- [email protected]
To unsubscribe send an email to [email protected]
https://mail.python.org/mailman3//lists/python-checkins.python.org
Member address: [email protected]
lmpx.com only provides a reader for public news (NNTP) servers. It is not affiliated with the servers or forums shown here and is not responsible for the content of articles, which is written by their respective authors.