@@ -215,11 +215,17 @@ def _cleanup_sockets(*sockets):
215215)
216216
217217
218- def _asyncio_in_subinterpreter ():
218+ def _asyncio_in_subinterpreter (ready , release ):
219+ """Park an asyncio task in a subinterpreter until released."""
219220 import asyncio
220221
222+ def wait (loop ):
223+ # Signal via the loop, so the task has already suspended
224+ loop .call_soon_threadsafe (ready .put , None )
225+ release .get ()
226+
221227 async def sub_worker ():
222- await asyncio .sleep ( 2 )
228+ await asyncio .to_thread ( wait , asyncio . get_running_loop () )
223229
224230 asyncio .run (sub_worker ())
225231
@@ -506,30 +512,36 @@ async def main():
506512 @requires_subinterpreters
507513 def test_all_awaited_by_covers_every_interpreter (self ):
508514 # gh-158880
515+ ready = interpreters .create_queue ()
516+ release = interpreters .create_queue ()
517+
509518 async def main_worker ():
510519 await asyncio .sleep (SHORT_TIMEOUT )
511520
512521 async def main ():
513522 with InterpreterPoolExecutor () as pool :
514- loop = asyncio .get_running_loop ()
515- loop .run_in_executor (pool , _asyncio_in_subinterpreter )
516- task = asyncio .create_task (main_worker (), name = "main_worker" )
517- self .addCleanup (task .cancel )
518- for _ in busy_retry (SHORT_TIMEOUT ):
523+ try :
524+ loop = asyncio .get_running_loop ()
525+ loop .run_in_executor (pool , _asyncio_in_subinterpreter ,
526+ ready , release )
527+ task = asyncio .create_task (main_worker (),
528+ name = "main_worker" )
529+ self .addCleanup (task .cancel )
519530 await asyncio .sleep (0 )
520- stacks = [
531+ ready .get (timeout = SHORT_TIMEOUT )
532+ return [
521533 [frame .funcname .rpartition ("." )[2 ]
522534 for frame in coro .call_stack ]
523535 for info in RemoteUnwinder (
524536 os .getpid ()).get_all_awaited_by ()
525537 for task in info .awaited_by
526538 for coro in task .coroutine_stack
527539 ]
528- if [ "sleep" , "sub_worker" ] in stacks :
529- return stacks
540+ finally :
541+ release . put ( None )
530542
531543 stacks = asyncio .run (main ())
532- self .assertIn (["sleep " , "sub_worker" ], stacks )
544+ self .assertIn (["to_thread " , "sub_worker" ], stacks )
533545 self .assertIn (["sleep" , "main_worker" ], stacks )
534546
535547 @skip_if_not_supported
0 commit comments