@@ -332,30 +332,39 @@ def clear(self):
332332
333333 @warnings_helper .ignore_fork_in_thread_deprecation_warnings ()
334334 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.
335+ # Workers must not write to the result queue on shutdown: the
336+ # executor manager thread does not read it while it joins them.
337+ # Fill the pipe once the executor manager thread stops reading
338+ # results, right before the workers are told to exit.
339339 self .executor .shutdown (wait = True )
340340
341341 data = b"a" * support .PIPE_MAX_SIZE
342- fillers = []
343- shutdown_workers = futures .process ._ExecutorManagerThread .shutdown_workers
344- def mock_shutdown_workers (self ):
342+ join_executor_internals = (
343+ futures .process ._ExecutorManagerThread ._join_executor_internals )
344+ def mock_join_executor_internals (self , broken = False ):
345345 # The put() blocks holding the result queue's write lock until
346- # the pipe is drained, so every exiting worker blocks behind it.
346+ # the pipe is drained, so any worker writing to the result queue
347+ # on exit blocks behind it.
347348 filler = threading .Thread (target = self .result_queue .put ,
348349 args = (data ,))
349350 filler .start ()
350- fillers .append (filler )
351- shutdown_workers (self )
351+ wlock = self .result_queue ._wlock
352+ if wlock is not None :
353+ # Wait for the filler to hold the write lock.
354+ while wlock .acquire (block = False ):
355+ wlock .release ()
356+ time .sleep (0.001 )
357+ join_executor_internals (self , broken )
358+ # Unblock the filler.
359+ self .result_queue .get ()
360+ filler .join ()
352361
353362 executor = self .executor_type (max_workers = 2 ,
354363 mp_context = self .get_context ())
355364 self .executor = executor # Allow clean up in fail_on_deadlock
356365 with unittest .mock .patch .object (futures .process ._ExecutorManagerThread ,
357- 'shutdown_workers ' ,
358- mock_shutdown_workers ):
366+ '_join_executor_internals ' ,
367+ mock_join_executor_internals ):
359368 self .assertEqual (list (executor .map (int , range (10 ))),
360369 list (range (10 )))
361370 shutdown = threading .Thread (target = executor .shutdown )
@@ -364,9 +373,6 @@ def mock_shutdown_workers(self):
364373 if shutdown .is_alive ():
365374 self ._fail_on_deadlock (executor )
366375
367- self .assertEqual (len (fillers ), 1 )
368- fillers [0 ].join ()
369-
370376
371377create_executor_tests (globals (), ExecutorDeadlockTest ,
372378 executor_mixins = (ProcessPoolForkMixin ,
0 commit comments