diff --git a/Lib/asyncio/locks.py b/Lib/asyncio/locks.py index c59ea15111bef2b..ec3abc55aa8af7a 100644 --- a/Lib/asyncio/locks.py +++ b/Lib/asyncio/locks.py @@ -489,6 +489,7 @@ def __init__(self, parties): self._parties = parties self._state = _BarrierState.FILLING self._count = 0 # count tasks in Barrier + self._index = self._parties - 1 # index number returned when barrier drains def __repr__(self): res = super().__repr__() @@ -515,13 +516,14 @@ async def wait(self): async with self._cond: await self._block() # Block while the barrier drains or resets. try: - index = self._count self._count += 1 - if index + 1 == self._parties: + if self._count == self._parties: # We release the barrier await self._release() else: await self._wait() + index = self._index + self._index = (self._index + 1) % self.parties return index finally: self._count -= 1 @@ -569,6 +571,7 @@ def _exit(self): if self._count == 0: if self._state in (_BarrierState.RESETTING, _BarrierState.DRAINING): self._state = _BarrierState.FILLING + self._index = self.parties - 1 self._cond.notify_all() async def reset(self): @@ -584,6 +587,7 @@ async def reset(self): self._state = _BarrierState.RESETTING else: self._state = _BarrierState.FILLING + self._index = self.parties - 1 self._cond.notify_all() async def abort(self): diff --git a/Lib/test/test_asyncio/test_locks.py b/Lib/test/test_asyncio/test_locks.py index 320a933e144813d..c6ab38a696a693d 100644 --- a/Lib/test/test_asyncio/test_locks.py +++ b/Lib/test/test_asyncio/test_locks.py @@ -1820,6 +1820,35 @@ async def coro(): self.assertEqual(barrier1.n_waiting, 0) + async def test_filling_tasks_cancel_one_index_still_unique(self): + # See gh-155233: a task cancelled while the barrier is still + # filling used to leave its index available for reuse. + # Index counter is now calculated on demand during the draining phasis. + self.N = 3 + barrier = asyncio.Barrier(self.N) + results = [] + + async def coro(): + i = await barrier.wait() + results.append(i) + + t1 = asyncio.create_task(coro()) + t2 = asyncio.create_task(coro()) + await asyncio.sleep(0) + self.assertEqual(barrier.n_waiting, 2) + + t1.cancel() + with self.assertRaises(asyncio.CancelledError): + await t1 + await asyncio.sleep(0) + self.assertEqual(barrier.n_waiting, 1) + + t3 = asyncio.create_task(coro()) + t4 = asyncio.create_task(coro()) + await asyncio.gather(t2, t3, t4) + + self.assertEqual(sorted(results), list(range(self.N))) + self.assertEqual(barrier.n_waiting, 0) if __name__ == '__main__': unittest.main() diff --git a/Misc/NEWS.d/next/Library/2026-10-07-10-59-27.gh-issue-155233.YGuQGK.rst b/Misc/NEWS.d/next/Library/2026-10-07-10-59-27.gh-issue-155233.YGuQGK.rst new file mode 100644 index 000000000000000..bc1e35d4116a85c --- /dev/null +++ b/Misc/NEWS.d/next/Library/2026-10-07-10-59-27.gh-issue-155233.YGuQGK.rst @@ -0,0 +1 @@ +Fix the returned value of the :meth:`asyncio.Barrier.wait`. This value is now calculated on demand during the *draining* phasis.