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
99 changes: 57 additions & 42 deletions Lib/test/_test_multiprocessing.py
Original file line number Diff line number Diff line change
Expand Up @@ -1220,8 +1220,10 @@ def test_put(self):
queue.put(5, False, None)
queue.put_nowait(6)

# the values may be in buffer but not yet in pipe so sleep a bit
time.sleep(DELTA)
# The values may be in the buffer but not yet in the pipe.
for _ in support.sleeping_retry(support.SHORT_TIMEOUT):
if not queue_empty(queue):
break

self.assertEqual(queue_empty(queue), False)
self.assertEqual(queue_full(queue, MAXSIZE), True)
Expand Down Expand Up @@ -1342,8 +1344,10 @@ def test_fork(self):
for i in range(10):
queue.put(i)

# wait to make sure thread starts before we fork a new process
time.sleep(DELTA)
# Wait for the feeder thread to write to the pipe before forking.
for _ in support.sleeping_retry(support.SHORT_TIMEOUT):
if not queue_empty(queue):
break

# fork process
p = self.Process(target=self._test_fork, args=(queue,))
Expand Down Expand Up @@ -1566,15 +1570,15 @@ def test_repr_lock(self):
args=(lock, l),
name=tname)
t.start()
time.sleep(0.1)
threading_helper.join_thread(t)
self.assertEqual(f'<Lock(owner=MainProcess|{tname})>', l[0])
lock.release()

t = threading.Thread(target=self._acquire,
args=(lock,),
name=tname)
t.start()
time.sleep(0.1)
threading_helper.join_thread(t)
self.assertEqual('<Lock(owner=SomeOtherThread)>', repr(lock))
lock.release()

Expand Down Expand Up @@ -3175,25 +3179,28 @@ def produce_args():
last_produced_task_arg.value = arg
yield arg

method = getattr(p, method_name)
it = method(functools.partial(sqr, wait=0.2), produce_args())
def wait_for_produced_args(threshold):
for _ in support.sleeping_retry(support.SHORT_TIMEOUT, error=False):
if last_produced_task_arg.value > threshold:
break
self.assertGreater(last_produced_task_arg.value, threshold)

time.sleep(0.2)
# `iterable` could've been advanced only `processes` times,
# but in fact it advances further (`> processes`) because of
# not waiting for workers or user code to catch up.
self.assertGreater(last_produced_task_arg.value, processes)
try:
method = getattr(p, method_name)
it = method(functools.partial(sqr, wait=0.2), produce_args())

next(it)
time.sleep(0.2)
self.assertGreater(last_produced_task_arg.value, processes + 1)
# Without a buffer limit, the iterable advances without waiting
# for workers or user code to catch up.
wait_for_produced_args(processes)

next(it)
time.sleep(0.2)
self.assertGreater(last_produced_task_arg.value, processes + 2)
next(it)
wait_for_produced_args(processes + 1)

p.terminate()
p.join()
next(it)
wait_for_produced_args(processes + 2)
finally:
p.terminate()
p.join()

@unittest.skipUnless(HAS_SHAREDCTYPES, 'needs sharedctypes')
@warnings_helper.ignore_fork_in_thread_deprecation_warnings()
Expand All @@ -3213,23 +3220,29 @@ def produce_args():
last_produced_task_arg.value = arg
yield arg

method = getattr(p, method_name)
it = method(functools.partial(sqr, wait=0.2), produce_args(),
buffersize=processes)
def wait_for_produced_args(expected):
# Consuming a result allows the task handler to advance the
# iterable, but does not wait for it to do so.
for _ in support.sleeping_retry(support.SHORT_TIMEOUT, error=False):
if last_produced_task_arg.value >= expected:
break
self.assertEqual(last_produced_task_arg.value, expected)

time.sleep(0.2)
self.assertEqual(last_produced_task_arg.value, processes)
try:
method = getattr(p, method_name)
it = method(functools.partial(sqr, wait=0.2), produce_args(),
buffersize=processes)

next(it)
time.sleep(0.2)
self.assertEqual(last_produced_task_arg.value, processes + 1)
wait_for_produced_args(processes)

next(it)
time.sleep(0.2)
self.assertEqual(last_produced_task_arg.value, processes + 2)
next(it)
wait_for_produced_args(processes + 1)

p.terminate()
p.join()
next(it)
wait_for_produced_args(processes + 2)
finally:
p.terminate()
p.join()

@warnings_helper.ignore_fork_in_thread_deprecation_warnings()
@support.subTests('method_name', ("imap", "imap_unordered"))
Expand Down Expand Up @@ -3435,8 +3448,11 @@ def test_release_task_refs(self):
self.pool.map(identity, objs)

del objs
time.sleep(DELTA) # let threaded cleanup code run
support.gc_collect() # For PyPy or other GCs.
for _ in support.sleeping_retry(support.SHORT_TIMEOUT, error=False):
support.gc_collect() # For PyPy or other GCs.
if (all(wr() is None for wr in refs)
and CountedObject.n_instances == 0):
break
self.assertEqual(set(wr() for wr in refs), {None})
# With a process pool, copies of the objects are returned, check
# they were released too.
Expand Down Expand Up @@ -4030,7 +4046,7 @@ def test_connection(self):
self.assertTimingAlmostEqual(poll.elapsed, TIMEOUT1)

conn.send(None)
time.sleep(.1)
self.assertTrue(conn.poll(support.SHORT_TIMEOUT))

self.assertEqual(poll(TIMEOUT1), True)
self.assertTimingAlmostEqual(poll.elapsed, 0)
Expand Down Expand Up @@ -4331,11 +4347,10 @@ def test_issue14725(self):
p = self.Process(target=self._test, args=(l.address,))
p.daemon = True
p.start()
time.sleep(1)
# On Windows the client process should by now have connected,
# written data and closed the pipe handle by now. This causes
# ConnectNamdedPipe() to fail with ERROR_NO_DATA. See Issue
# 14725.
join_process(p)
# On Windows the client process has now connected, written data
# and closed the pipe handle. This causes ConnectNamedPipe()
# to fail with ERROR_NO_DATA. See Issue 14725.
conn = l.accept()
self.assertEqual(conn.recv(), 'hello')
conn.close()
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
Replace fixed sleeps used for synchronization in multiprocessing tests with
bounded waits to avoid failures on slow or overloaded systems. Ensure pool
cleanup runs even when assertions fail in the imap buffer tests.
Loading