diff --git a/Lib/multiprocessing/synchronize.py b/Lib/multiprocessing/synchronize.py index 9188114ae284c7..e3ecaa5403b12b 100644 --- a/Lib/multiprocessing/synchronize.py +++ b/Lib/multiprocessing/synchronize.py @@ -125,6 +125,23 @@ def _make_name(): return '%s-%s' % (process.current_process()._config['semprefix'], next(SemLock._rand)) + def _get_procname_and_count(self): + try: + if self._semlock._is_mine(): + name = process.current_process().name + if threading.current_thread().name != 'MainThread': + name += '|' + threading.current_thread().name + count = self._semlock._count() + elif not self._semlock._is_zero(): + name, count = 'None', 0 + elif self._semlock._count() > 0: + name, count = 'SomeOtherThread', 'nonzero' + else: + name, count = 'SomeOtherProcess', 'nonzero' + except Exception: + name, count = 'unknown', 'unknown' + return name, count + # # Semaphore # @@ -143,11 +160,12 @@ def get_value(self): return self._semlock._get_value() def __repr__(self): + res = super().__repr__() try: value = self.get_value() except Exception: value = 'unknown' - return '<%s(value=%s)>' % (self.__class__.__name__, value) + return f'<{res[1:-1]} (value={value})>' # # Bounded semaphore @@ -159,12 +177,13 @@ def __init__(self, value=1, *, ctx): SemLock.__init__(self, SEMAPHORE, value, value, ctx=ctx) def __repr__(self): + res = object.__repr__(self) try: value = self.get_value() except Exception: value = 'unknown' - return '<%s(value=%s, maxvalue=%s)>' % \ - (self.__class__.__name__, value, self._semlock.maxvalue) + return f'<{res[1:-1]} (value={value}, ' \ + f'maxvalue={self._semlock.maxvalue})>' # # Non-recursive lock @@ -176,20 +195,9 @@ def __init__(self, *, ctx): SemLock.__init__(self, SEMAPHORE, 1, 1, ctx=ctx) def __repr__(self): - try: - if self._semlock._is_mine(): - name = process.current_process().name - if threading.current_thread().name != 'MainThread': - name += '|' + threading.current_thread().name - elif not self._semlock._is_zero(): - name = 'None' - elif self._semlock._count() > 0: - name = 'SomeOtherThread' - else: - name = 'SomeOtherProcess' - except Exception: - name = 'unknown' - return '<%s(owner=%s)>' % (self.__class__.__name__, name) + res = super().__repr__() + name, _ = self._get_procname_and_count() + return f'<{res[1:-1]} (owner={name})>' # # Recursive lock @@ -201,21 +209,9 @@ def __init__(self, *, ctx): SemLock.__init__(self, RECURSIVE_MUTEX, 1, 1, ctx=ctx) def __repr__(self): - try: - if self._semlock._is_mine(): - name = process.current_process().name - if threading.current_thread().name != 'MainThread': - name += '|' + threading.current_thread().name - count = self._semlock._count() - elif not self._semlock._is_zero(): - name, count = 'None', 0 - elif self._semlock._count() > 0: - name, count = 'SomeOtherThread', 'nonzero' - else: - name, count = 'SomeOtherProcess', 'nonzero' - except Exception: - name, count = 'unknown', 'unknown' - return '<%s(%s, %s)>' % (self.__class__.__name__, name, count) + res = super().__repr__() + name, count = self._get_procname_and_count() + return f'<{res[1:-1]} (owner={name}, count={count})>' # # Condition variable @@ -251,12 +247,15 @@ def _make_methods(self): self.release = self._lock.release def __repr__(self): + res = super().__repr__() try: num_waiters = (self._sleeping_count.get_value() - self._woken_count.get_value()) except Exception: num_waiters = 'unknown' - return '<%s(%s, %s)>' % (self.__class__.__name__, self._lock, num_waiters) + + lock_repr, _ = self._lock._get_procname_and_count() + return f'<{res[1:-1]} (lock={lock_repr}, waiters={num_waiters})>' def wait(self, timeout=None): assert self._lock._semlock._is_mine(), \ @@ -368,8 +367,10 @@ def wait(self, timeout=None): return False def __repr__(self): + res = super().__repr__() set_status = 'set' if self.is_set() else 'unset' - return f"<{type(self).__qualname__} at {id(self):#x} {set_status}>" + return f'<{res[1:-1]} ({set_status})>' + # # Barrier # @@ -409,3 +410,9 @@ def _count(self): @_count.setter def _count(self, value): self._array[1] = value + + def __repr__(self): + res = object.__repr__(self) + if self.broken: + return f'<{res[1:-1]} (broken)>' + return f'<{res[1:-1]} (waiters={self.n_waiting}/{self.parties})>' diff --git a/Lib/test/_test_multiprocessing.py b/Lib/test/_test_multiprocessing.py index 36e0880bc08818..58415ce02b623d 100644 --- a/Lib/test/_test_multiprocessing.py +++ b/Lib/test/_test_multiprocessing.py @@ -138,8 +138,11 @@ def _resource_unlink(name, rtype): WAIT_ACTIVE_CHILDREN_TIMEOUT = 5.0 -HAVE_GETVALUE = not getattr(_multiprocessing, - 'HAVE_BROKEN_SEM_GETVALUE', False) +try: + HAVE_GETVALUE = not getattr(_multiprocessing,'flags') \ + .get('HAVE_BROKEN_SEM_GETVALUE', False) +except: + HAVE_GETVALUE = True WIN32 = (sys.platform == "win32") @@ -1545,7 +1548,7 @@ def _acquire(lock, l=None): def _acquire_event(lock, event): lock.acquire() event.set() - time.sleep(1.0) + time.sleep(0.1) @warnings_helper.ignore_fork_in_thread_deprecation_warnings() def test_repr_lock(self): @@ -1553,10 +1556,12 @@ def test_repr_lock(self): self.skipTest('test not appropriate for {}'.format(self.TYPE)) lock = self.Lock() - self.assertEqual(f'', repr(lock)) + self.assertRegex(repr(lock), r"<*.Lock object at .* \(owner=None\)>") lock.acquire() - self.assertEqual(f'', repr(lock)) + self.assertRegex(repr(lock), + r"<*.Lock object at .* " + r"\(owner=MainProcess\)>") lock.release() tname = 'T1' @@ -1566,16 +1571,21 @@ def test_repr_lock(self): name=tname) t.start() time.sleep(0.1) - self.assertEqual(f'', l[0]) + self.assertRegex(repr(l[0]), "<*.Lock object at .* " + f"\\(owner=MainProcess\\|{tname}\\)>") lock.release() + t.join() t = threading.Thread(target=self._acquire, args=(lock,), name=tname) t.start() time.sleep(0.1) - self.assertEqual('', repr(lock)) + self.assertRegex(repr(lock), + r"<*.Lock object at .* " + r"\(owner=SomeOtherThread\)>") lock.release() + t.join() pname = 'P1' l = multiprocessing.Manager().list() @@ -1584,7 +1594,7 @@ def test_repr_lock(self): name=pname) p.start() p.join() - self.assertEqual(f'', l[0]) + self.assertRegex(l[0], f"<*.Lock object at .* \\(owner={pname}\\)>") lock = self.Lock() event = self.Event() @@ -1593,8 +1603,10 @@ def test_repr_lock(self): name='P2') p.start() event.wait() - self.assertEqual(f'', repr(lock)) - p.terminate() + self.assertRegex(repr(lock), + f"<*.Lock object at .* " + f"\\(owner=SomeOtherProcess\\)>") + p.join() def test_lock(self): lock = self.Lock() @@ -1629,61 +1641,77 @@ def test_lock_locked_2processes(self): p.join() @staticmethod - def _acquire_release(lock, timeout, l=None, n=1): + def _acquire_release(rlock, timeout, l=None, n=1): for _ in range(n): - lock.acquire() + rlock.acquire() if l is not None: - l.append(repr(lock)) + l.append(repr(rlock)) time.sleep(timeout) for _ in range(n): - lock.release() + rlock.release() @warnings_helper.ignore_fork_in_thread_deprecation_warnings() def test_repr_rlock(self): if self.TYPE != 'processes': self.skipTest('test not appropriate for {}'.format(self.TYPE)) - lock = self.RLock() - self.assertEqual('', repr(lock)) + rlock = self.RLock() + self.assertRegex(repr(rlock), + r"<*.RLock object at .* " + r"\(owner=None, count=0\)>") n = 3 for _ in range(n): - lock.acquire() - self.assertEqual(f'', repr(lock)) + rlock.acquire() + self.assertRegex(repr(rlock), + f"<*.RLock object at .* " + f"\\(owner=MainProcess, count={n}\\)>") + for _ in range(n): - lock.release() + rlock.release() t, l = [], [] for i in range(n): t.append(threading.Thread(target=self._acquire_release, - args=(lock, 0.1, l, i+1), + args=(rlock, 0.1, l, i+1), name=f'T{i+1}')) t[-1].start() for t_ in t: t_.join() for i in range(n): - self.assertIn(f'', l) + if i < len(l): + self.assertRegex(repr(l[i]), + f"<*.RLock object at .* " + f"\\(owner=MainProcess|T{i+1}, " + f"count=nonzero\\)>") rlock = self.RLock() t = threading.Thread(target=rlock.acquire) t.start() t.join() - self.assertEqual('', repr(rlock)) + self.assertRegex(repr(rlock), + r"<*.RLock object at .* " + r"\(owner=SomeOtherThread, count=nonzero\)>") pname = 'P1' + rlock = self.RLock() l = multiprocessing.Manager().list() p = self.Process(target=self._acquire_release, - args=(lock, 0.1, l), + args=(rlock, 0.1, l), name=pname) p.start() p.join() - self.assertEqual(f'', l[0]) + self.assertRegex(repr(l[0]), + f"<*.RLock object at .* " + f"\\(owner={pname}, count=1\\)>") rlock = self.RLock() p = self.Process(target=self._acquire, args=(rlock,)) p.start() p.join() - self.assertEqual('', repr(rlock)) + self.assertRegex(repr(rlock), + r"<*.RLock object at .* " + r"\(owner=SomeOtherProcess, count=nonzero\)>") def test_rlock(self): lock = self.RLock() @@ -1752,9 +1780,36 @@ def test_bounded_semaphore(self): sem = self.BoundedSemaphore(2) self._test_semaphore(sem) # Currently fails on OS/X - #if HAVE_GETVALUE: - # self.assertRaises(ValueError, sem.release) - # self.assertReturnsIfImplemented(2, get_value, sem) + if HAVE_GETVALUE: + self.assertRaises(ValueError, sem.release) + self.assertReturnsIfImplemented(2, get_value, sem) + + @unittest.skipUnless(HAVE_GETVALUE, 'needs sem_getvalue') + def test_repr_semaphore(self): + if self.TYPE != 'processes': + self.skipTest('test not appropriate for {}'.format(self.TYPE)) + n = 5 + sem = self.Semaphore(n) + self.assertRegex(repr(sem), + f"<*.Semaphore object at .* \\(value={n}\\)>") + sem.acquire() + self.assertRegex(repr(sem), + f"<*.Semaphore object at .* \\(value={n-1}\\)>") + + @unittest.skipUnless(HAVE_GETVALUE, 'needs sem_getvalue') + def test_repr_boundedsemaphore(self): + if self.TYPE != 'processes': + self.skipTest('test not appropriate for {}'.format(self.TYPE)) + n = 5 + bsem = self.BoundedSemaphore(n) + self.assertRegex(repr(bsem), + f"<*.BoundedSemaphore object at .* " + f"\\(value={n}, maxvalue={n}\\)>") + for _ in range(n): + bsem.acquire() + self.assertRegex(repr(bsem), + f"<*.BoundedSemaphore object at .* " + f"\\(value=0, maxvalue={n}\\)>") def test_timeout(self): if self.TYPE != 'processes': @@ -1810,6 +1865,57 @@ def check_invariant(self, cond): except NotImplementedError: pass + @warnings_helper.ignore_fork_in_thread_deprecation_warnings() + def test_repr_condition(self): + cond = self.Condition() + if self.TYPE == 'processes': + # on macOS, waiters can be 'unknown'. + self.assertRegex(repr(cond), + r"<*.Condition object at .* " + r"\(lock=None, waiters=.*\)>") + cond.acquire() + self.assertRegex(repr(cond), + r"<*.Condition object at .* " + r"\(lock=MainProcess, waiters=.*\)>") + cond.release() + self.assertRegex(repr(cond), + r"<*.Condition object at .* " + r"\(lock=None, waiters=.*\)>") + + elif self.TYPE == 'managers': + self.assertRegex(repr(cond), + r"ConditionProxy object, typeid 'Condition' at .*>") + cond.acquire() + print(cond) + + @staticmethod + def _cond_wait(cond): + with cond: + cond.wait() + + @unittest.skipUnless(HAVE_GETVALUE, 'needs sem_getvalue') + @warnings_helper.ignore_fork_in_thread_deprecation_warnings() + def test_repr_condition_waiters(self): + cond = self.Condition() + if self.TYPE == 'processes': + p = self.Process(target=self._cond_wait, args=(cond,)) + p.start() + # No need to check _woken_count attr. + while cond._sleeping_count.get_value() == 0: + _wait() + self.assertRegex(repr(cond), + r"<*.Condition object at .* " + r"\(lock=None, waiters=1\)>") + with cond: + self.assertRegex(repr(cond), + r"<*.Condition object at .* " + r"\(lock=MainProcess, waiters=1\)>") + cond.notify(1) + p.join() + self.assertRegex(repr(cond), + r"<*.Condition object at .* " + r"\(lock=None, waiters=0\)>") + @warnings_helper.ignore_fork_in_thread_deprecation_warnings() def test_notify(self): cond = self.Condition() @@ -2131,18 +2237,20 @@ def test_event(self): self.assertEqual(wait(), True) p.join() - def test_repr(self) -> None: + def test_repr_event(self) -> None: event = self.Event() if self.TYPE == 'processes': - self.assertRegex(repr(event), r"") + self.assertRegex(repr(event), r"<*.Event object at .* \(unset\)>") event.set() - self.assertRegex(repr(event), r"") + self.assertRegex(repr(event), r"<*.Event object at .* \(set\)>") event.clear() - self.assertRegex(repr(event), r"") + self.assertRegex(repr(event), r"<*.Event object at .* \(unset\)>") elif self.TYPE == 'manager': - self.assertRegex(repr(event), r"") + self.barrier.abort() + self.assertRegex(repr(self.barrier), + f"<*.Barrier object at .* \\(broken\\)>") + self.barrier.reset() + ps = [] + for i in range(self.N-1): + p = self.Process(target=self._barrier_wait, + args=(self.barrier,)) + p.start() + ps.append(p) + # wait for at least one processes waits + while self.barrier.n_waiting == 0: + _wait() + self.assertRegex(repr(self.barrier), + f"<*.Barrier object at .* " + f"\\(waiters={list(range(1, self.N))}/{self.N}\\)>") + self.barrier.wait() + for p in ps: + p.join() + self.assertRegex(repr(self.barrier), + f"<*.Barrier object at .* \\(waiters=0/{self.N}\\)>") + elif self.TYPE == 'manager': + self.assertRegex(repr(self.barrier), + r"