[FIX] asyncio : rendre --max_process lançable sur Python 3.10 et plus
Le drapeau `--max_process` de deux scripts passe par ce pool, et il ne tournait plus du tout. Deux retraits d'API le traversaient : `loop=` a quitté `asyncio.wait` en 3.10, où le passer lève un TypeError, et `asyncio.get_event_loop()` lève hors d'une loop en marche depuis 3.14, où il ne faisait qu'avertir — donc la classe n'était même plus instanciable, alors que l'aide annonce toujours l'option. La loop n'est plus créée à l'instanciation mais à l'exécution, et `close` ne ferme que celle que la classe a ouverte : une loop reçue en argument appartient à l'appelant, qui compte encore dessus. Vérifié : 7 tests sur de vraies coroutines, et les deux appels d'origine lèvent toujours à part. --- EN --- The `--max_process` flag of two scripts goes through this pool, and it no longer ran at all. Two API removals crossed it: `loop=` left `asyncio.wait` in 3.10, where passing it raises a TypeError, and `asyncio.get_event_loop()` raises outside a running loop since 3.14, where it merely warned before — so the class was not even constructible, while the help still advertises the option. The loop is no longer created at construction but at run time, and `close` only closes the one the class opened: a loop received as an argument belongs to the caller, who still counts on it. Checked: 7 tests on real coroutines, and both original calls still raise on their own. Assisted-by: Claude Opus 5
This commit is contained in:
parent
d335293d85
commit
ea6391a781
2 changed files with 137 additions and 4 deletions
|
|
@ -155,15 +155,30 @@ class AsyncioPool:
|
|||
"""
|
||||
@param loop: asyncio loop
|
||||
@param concurrency: Maximum number of concurrently running tasks
|
||||
|
||||
La loop n'est PAS créée ici. `asyncio.get_event_loop()` lève hors
|
||||
d'une loop en marche depuis Python 3.14 — il ne faisait qu'avertir
|
||||
avant — donc la construire à l'instanciation rendait la classe
|
||||
impossible à instancier, et le drapeau `--max_process` inutilisable.
|
||||
Elle est donc créée à l'exécution, et fermée par `close`, qui ne
|
||||
ferme que ce que la classe a ouvert.
|
||||
"""
|
||||
self._loop = loop or asyncio.get_event_loop()
|
||||
self._loop = loop
|
||||
self._loop_est_notre = False
|
||||
self._concurrency = concurrency
|
||||
self._coros = deque([]) # All coroutines queued for execution
|
||||
self._futures = [] # All currently running coroutines
|
||||
self._lst_result = []
|
||||
|
||||
def close(self):
|
||||
"""Ferme la loop, si c'est celle que la classe a créée.
|
||||
|
||||
Une loop reçue en argument appartient à l'appelant : la fermer sous
|
||||
lui casserait tout ce qu'il compte encore y faire tourner."""
|
||||
if self._loop is not None and self._loop_est_notre:
|
||||
self._loop.close()
|
||||
self._loop = None
|
||||
self._loop_est_notre = False
|
||||
|
||||
def add_coro(self, coro):
|
||||
"""
|
||||
|
|
@ -173,6 +188,10 @@ class AsyncioPool:
|
|||
self.print_status()
|
||||
|
||||
def run_until_complete(self):
|
||||
if self._loop is None:
|
||||
self._loop = asyncio.new_event_loop()
|
||||
self._loop_est_notre = True
|
||||
asyncio.set_event_loop(self._loop)
|
||||
self._loop.run_until_complete(self._wait_for_futures())
|
||||
return self._lst_result
|
||||
|
||||
|
|
@ -187,16 +206,19 @@ class AsyncioPool:
|
|||
num_to_start = min(num_to_start, len(self._coros))
|
||||
for _ in range(num_to_start):
|
||||
coro = self._coros.popleft()
|
||||
future = asyncio.ensure_future(coro, loop=self._loop)
|
||||
# Sans « loop= » : l'appel a lieu DEPUIS la loop en marche, qui
|
||||
# est donc celle que ensure_future prend d'elle-même.
|
||||
future = asyncio.ensure_future(coro)
|
||||
self._futures.append(future)
|
||||
self.print_status()
|
||||
|
||||
async def _wait_for_futures(self):
|
||||
while len(self._coros) > 0 or len(self._futures) > 0:
|
||||
self._start_futures()
|
||||
# « loop= » a été RETIRÉ de asyncio.wait en Python 3.10 : le
|
||||
# passer lève un TypeError, et c'est la loop en marche qui sert.
|
||||
futures_completed, futures_pending = await asyncio.wait(
|
||||
self._futures,
|
||||
loop=self._loop,
|
||||
return_when=asyncio.FIRST_COMPLETED,
|
||||
)
|
||||
|
||||
|
|
|
|||
111
test/test_lib_asyncio_pool.py
Normal file
111
test/test_lib_asyncio_pool.py
Normal file
|
|
@ -0,0 +1,111 @@
|
|||
#!/usr/bin/env python3
|
||||
# © 2026 TechnoLibre (http://www.technolibre.ca)
|
||||
# License AGPL-3.0 or later (http://www.gnu.org/licenses/agpl)
|
||||
"""Le pool de coroutines, et les deux retraits d'asyncio qui le tuaient.
|
||||
|
||||
Le drapeau `--max_process` de `run_parallel_test.py` et de
|
||||
`show_evolution_module.py` passe par cette classe, et elle ne pouvait plus
|
||||
tourner du tout : `asyncio.get_event_loop()` LÈVE hors d'une loop en marche
|
||||
depuis Python 3.14, ce qui rendait l'instanciation impossible, et le mot-clé
|
||||
`loop=` a été RETIRÉ d'`asyncio.wait` en Python 3.10, ce qui lève un
|
||||
TypeError. Deux échecs à des étages différents, tous deux sur le chemin d'un
|
||||
drapeau que l'aide annonce encore.
|
||||
|
||||
Ces tests exercent le pool pour de vrai — de vraies coroutines, une vraie
|
||||
loop, un vrai résultat — parce que c'est la seule façon de prouver qu'un
|
||||
retrait d'API ne le traverse plus. Ils n'affirment rien sur l'ordre des
|
||||
résultats : le pool rend ce qui finit, dans l'ordre où cela finit.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import io
|
||||
import unittest
|
||||
from contextlib import redirect_stdout
|
||||
|
||||
from script.lib_asyncio import AsyncioPool
|
||||
|
||||
|
||||
async def rendre(valeur, delai=0):
|
||||
if delai:
|
||||
await asyncio.sleep(delai)
|
||||
return valeur
|
||||
|
||||
|
||||
class TestLePoolTourne(unittest.TestCase):
|
||||
"""Le pool s'instancie, tourne et se ferme sans loop fournie."""
|
||||
|
||||
def _lancer(self, pool):
|
||||
"""Exécute en silence : le pool imprime son état à chaque ajout."""
|
||||
with redirect_stdout(io.StringIO()):
|
||||
return pool.run_until_complete()
|
||||
|
||||
def tearDown(self):
|
||||
asyncio.set_event_loop(None)
|
||||
|
||||
def test_it_can_be_built_without_a_loop(self):
|
||||
"""Le cas qui levait : construire hors de toute loop en marche."""
|
||||
pool = AsyncioPool(2)
|
||||
self.assertIsNone(pool._loop)
|
||||
|
||||
def test_it_runs_more_coros_than_its_concurrency(self):
|
||||
pool = AsyncioPool(2)
|
||||
with redirect_stdout(io.StringIO()):
|
||||
for valeur in range(5):
|
||||
pool.add_coro(rendre(valeur))
|
||||
resultats = self._lancer(pool)
|
||||
pool.close()
|
||||
self.assertEqual(sorted(resultats), [0, 1, 2, 3, 4])
|
||||
|
||||
def test_a_single_coro(self):
|
||||
pool = AsyncioPool(4)
|
||||
with redirect_stdout(io.StringIO()):
|
||||
pool.add_coro(rendre("seule"))
|
||||
resultats = self._lancer(pool)
|
||||
pool.close()
|
||||
self.assertEqual(resultats, ["seule"])
|
||||
|
||||
def test_staggered_coros_all_come_back(self):
|
||||
"""Des durées inégales : c'est là que FIRST_COMPLETED est exercé."""
|
||||
pool = AsyncioPool(2)
|
||||
with redirect_stdout(io.StringIO()):
|
||||
for valeur, delai in ((1, 0.03), (2, 0.01), (3, 0.02), (4, 0)):
|
||||
pool.add_coro(rendre(valeur, delai))
|
||||
resultats = self._lancer(pool)
|
||||
pool.close()
|
||||
self.assertEqual(sorted(resultats), [1, 2, 3, 4])
|
||||
|
||||
|
||||
class TestLaFermeture(unittest.TestCase):
|
||||
"""`close` ne ferme que ce que la classe a ouvert."""
|
||||
|
||||
def tearDown(self):
|
||||
asyncio.set_event_loop(None)
|
||||
|
||||
def test_close_before_any_run_is_harmless(self):
|
||||
AsyncioPool(2).close()
|
||||
|
||||
def test_close_closes_the_loop_it_created(self):
|
||||
pool = AsyncioPool(1)
|
||||
with redirect_stdout(io.StringIO()):
|
||||
pool.add_coro(rendre(1))
|
||||
pool.run_until_complete()
|
||||
boucle = pool._loop
|
||||
pool.close()
|
||||
self.assertTrue(boucle.is_closed())
|
||||
|
||||
def test_a_borrowed_loop_is_left_open(self):
|
||||
"""Une loop reçue appartient à l'appelant, qui compte encore dessus."""
|
||||
boucle = asyncio.new_event_loop()
|
||||
try:
|
||||
pool = AsyncioPool(1, loop=boucle)
|
||||
with redirect_stdout(io.StringIO()):
|
||||
pool.add_coro(rendre(1))
|
||||
pool.run_until_complete()
|
||||
pool.close()
|
||||
self.assertFalse(boucle.is_closed())
|
||||
finally:
|
||||
boucle.close()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Loading…
Reference in a new issue