From 03b736b04b0a7df293eaf846b22330664a9e5323 Mon Sep 17 00:00:00 2001 From: Timofey Ivankov Date: Tue, 6 Oct 2026 00:26:52 +0300 Subject: [PATCH 1/5] gh-158880: Fix asyncio ps showing tasks from only one interpreter --- Lib/test/test_external_inspection.py | 37 ++++++++++ ...-10-05-23-40-53.gh-issue-158880.8imLqO.rst | 2 + Modules/_remote_debugging/_remote_debugging.h | 1 + Modules/_remote_debugging/module.c | 72 ++++++++++++++----- Modules/_remote_debugging/threads.c | 3 +- 5 files changed, 97 insertions(+), 18 deletions(-) create mode 100644 Misc/NEWS.d/next/Library/2026-10-05-23-40-53.gh-issue-158880.8imLqO.rst diff --git a/Lib/test/test_external_inspection.py b/Lib/test/test_external_inspection.py index 9fff8e8bff91f7..960bb31c9e4e69 100644 --- a/Lib/test/test_external_inspection.py +++ b/Lib/test/test_external_inspection.py @@ -46,6 +46,7 @@ try: from concurrent import interpreters + from concurrent.futures import InterpreterPoolExecutor except ImportError: interpreters = None @@ -214,6 +215,15 @@ def _cleanup_sockets(*sockets): ) +def _asyncio_in_subinterpreter(): + import asyncio + + async def sub_worker(): + await asyncio.sleep(2) + + asyncio.run(sub_worker()) + + def requires_subinterpreters(meth): """Decorator to skip a test if subinterpreters are not supported.""" return unittest.skipIf(interpreters is None, "subinterpreters required")( @@ -492,6 +502,33 @@ async def main(): self.assertIn(main_name, names) self.assertEqual([len(n) for n in names if n.startswith("x")], [255]) + @skip_if_not_supported + @requires_subinterpreters + def test_all_awaited_by_covers_every_interpreter(self): + # gh-158880 + async def main_worker(): + await asyncio.sleep(2) + + async def main(): + with InterpreterPoolExecutor() as pool: + loop = asyncio.get_running_loop() + loop.run_in_executor(pool, _asyncio_in_subinterpreter) + task = asyncio.create_task(main_worker(), name="main_worker") + self.addCleanup(task.cancel) + await asyncio.sleep(1) + return [ + [frame.funcname.rpartition(".")[2] + for frame in coro.call_stack] + for info in RemoteUnwinder( + os.getpid()).get_all_awaited_by() + for task in info.awaited_by + for coro in task.coroutine_stack + ] + + stacks = asyncio.run(main()) + self.assertIn(["sleep", "sub_worker"], stacks) + self.assertIn(["sleep", "main_worker"], stacks) + @skip_if_not_supported def test_recursive_coroutine_stack_is_not_truncated(self): # gh-158522 diff --git a/Misc/NEWS.d/next/Library/2026-10-05-23-40-53.gh-issue-158880.8imLqO.rst b/Misc/NEWS.d/next/Library/2026-10-05-23-40-53.gh-issue-158880.8imLqO.rst new file mode 100644 index 00000000000000..5ab23d2f648610 --- /dev/null +++ b/Misc/NEWS.d/next/Library/2026-10-05-23-40-53.gh-issue-158880.8imLqO.rst @@ -0,0 +1,2 @@ +Fix ``python -m asyncio ps`` and ``pstree`` showing tasks from only one +interpreter. Patch by Timofei Ivankov. diff --git a/Modules/_remote_debugging/_remote_debugging.h b/Modules/_remote_debugging/_remote_debugging.h index 79fa4a92e745bd..8fc8fb05830228 100644 --- a/Modules/_remote_debugging/_remote_debugging.h +++ b/Modules/_remote_debugging/_remote_debugging.h @@ -663,6 +663,7 @@ extern int collect_frames_with_cache( extern int iterate_threads( RemoteUnwinderObject *unwinder, + uintptr_t interpreter_addr, thread_processor_func processor, void *context ); diff --git a/Modules/_remote_debugging/module.c b/Modules/_remote_debugging/module.c index 6fa277d1065ed9..e6056439a11746 100644 --- a/Modules/_remote_debugging/module.c +++ b/Modules/_remote_debugging/module.c @@ -958,7 +958,9 @@ _remote_debugging_RemoteUnwinder_get_all_awaited_by_impl(RemoteUnwinderObject *s if (ensure_async_debug_offsets(self) < 0) { return NULL; } - if (refresh_generation_caches_for_interpreter(self, self->interpreter_addr) < 0) { + PyObject *seen = PySet_New(NULL); + if (seen == NULL) { + set_exception_cause(self, PyExc_MemoryError, "Failed to create interpreter set"); return NULL; } @@ -968,30 +970,65 @@ _remote_debugging_RemoteUnwinder_get_all_awaited_by_impl(RemoteUnwinderObject *s goto result_err; } - // Process all threads - if (iterate_threads(self, process_thread_for_awaited_by, result) < 0) { - goto result_err; - } + // gh-158880: Tasks live in every interpreter, not only the one at the list head + for (uintptr_t interp = self->interpreter_addr; interp != 0; ) { + PyObject *addr = PyLong_FromUnsignedLongLong(interp); + if (addr == NULL) { + set_exception_cause(self, PyExc_MemoryError, "Failed to create interpreter address"); + goto result_err; + } + Py_ssize_t seen_count = PySet_GET_SIZE(seen); + int marked = PySet_Add(seen, addr); + Py_DECREF(addr); + if (marked < 0) { + set_exception_cause(self, PyExc_RuntimeError, "Failed to mark interpreter as seen"); + goto result_err; + } + if (PySet_GET_SIZE(seen) == seen_count) { + // already walked + break; + } - uintptr_t head_addr = self->interpreter_addr - + (uintptr_t)self->async_debug_offsets.asyncio_interpreter_state.asyncio_tasks_head; + if (refresh_generation_caches_for_interpreter(self, interp) < 0) { + goto result_err; + } - // On top of a per-thread task lists used by default by asyncio to avoid - // contention, there is also a fallback per-interpreter list of tasks; - // any tasks still pending when a thread is destroyed will be moved to the - // per-interpreter task list. It's unlikely we'll find anything here, but - // interesting for debugging. - if (append_awaited_by(self, 0, head_addr, result)) - { - set_exception_cause(self, PyExc_RuntimeError, "Failed to append interpreter awaited_by in get_all_awaited_by"); - goto result_err; + // Process all threads + if (iterate_threads(self, interp, process_thread_for_awaited_by, result) < 0) { + goto result_err; + } + + uintptr_t head_addr = interp + + (uintptr_t)self->async_debug_offsets.asyncio_interpreter_state.asyncio_tasks_head; + + // On top of a per-thread task lists used by default by asyncio to avoid + // contention, there is also a fallback per-interpreter list of tasks; + // any tasks still pending when a thread is destroyed will be moved to + // the per-interpreter task list. It's unlikely we'll find anything + // here, but interesting for debugging. + if (append_awaited_by(self, 0, head_addr, result)) + { + set_exception_cause(self, PyExc_RuntimeError, "Failed to append interpreter awaited_by in get_all_awaited_by"); + goto result_err; + } + + if (_Py_RemoteDebug_PagedReadRemoteMemory( + &self->handle, + interp + (uintptr_t)self->debug_offsets.interpreter_state.next, + sizeof(void*), + &interp) < 0) { + set_exception_cause(self, PyExc_RuntimeError, "Failed to read next interpreter address"); + goto result_err; + } } _Py_RemoteDebug_ClearCache(&self->handle); + Py_DECREF(seen); return result; result_err: _Py_RemoteDebug_ClearCache(&self->handle); + Py_DECREF(seen); Py_XDECREF(result); return NULL; } @@ -1064,7 +1101,8 @@ _remote_debugging_RemoteUnwinder_get_async_stack_trace_impl(RemoteUnwinderObject } // Process all threads - if (iterate_threads(self, process_thread_for_async_stack_trace, result) < 0) { + if (iterate_threads(self, self->interpreter_addr, + process_thread_for_async_stack_trace, result) < 0) { goto result_err; } diff --git a/Modules/_remote_debugging/threads.c b/Modules/_remote_debugging/threads.c index 198134fe6cfbea..49d1e7a36653cc 100644 --- a/Modules/_remote_debugging/threads.c +++ b/Modules/_remote_debugging/threads.c @@ -27,6 +27,7 @@ int iterate_threads( RemoteUnwinderObject *unwinder, + uintptr_t interpreter_addr, thread_processor_func processor, void *context ) { @@ -37,7 +38,7 @@ iterate_threads( if (0 > _Py_RemoteDebug_PagedReadRemoteMemory( &unwinder->handle, - unwinder->interpreter_addr + (uintptr_t)unwinder->debug_offsets.interpreter_state.threads_head, + interpreter_addr + (uintptr_t)unwinder->debug_offsets.interpreter_state.threads_head, sizeof(void*), &thread_state_addr)) { From 51482c0bffa2d280f92a870498cc99806b8f49ca Mon Sep 17 00:00:00 2001 From: Timofey Ivankov Date: Tue, 6 Oct 2026 17:56:10 +0300 Subject: [PATCH 2/5] remove potential flake in test --- Lib/test/test_external_inspection.py | 21 ++++++++++++--------- 1 file changed, 12 insertions(+), 9 deletions(-) diff --git a/Lib/test/test_external_inspection.py b/Lib/test/test_external_inspection.py index 960bb31c9e4e69..c6332ec311fac1 100644 --- a/Lib/test/test_external_inspection.py +++ b/Lib/test/test_external_inspection.py @@ -515,15 +515,18 @@ async def main(): loop.run_in_executor(pool, _asyncio_in_subinterpreter) task = asyncio.create_task(main_worker(), name="main_worker") self.addCleanup(task.cancel) - await asyncio.sleep(1) - return [ - [frame.funcname.rpartition(".")[2] - for frame in coro.call_stack] - for info in RemoteUnwinder( - os.getpid()).get_all_awaited_by() - for task in info.awaited_by - for coro in task.coroutine_stack - ] + for _ in busy_retry(SHORT_TIMEOUT): + await asyncio.sleep(0) + stacks = [ + [frame.funcname.rpartition(".")[2] + for frame in coro.call_stack] + for info in RemoteUnwinder( + os.getpid()).get_all_awaited_by() + for task in info.awaited_by + for coro in task.coroutine_stack + ] + if ["sleep", "sub_worker"] in stacks: + return stacks stacks = asyncio.run(main()) self.assertIn(["sleep", "sub_worker"], stacks) From 7c205da0e47a5f800245c44c33269908605e28bf Mon Sep 17 00:00:00 2001 From: Timofey Ivankov Date: Tue, 6 Oct 2026 19:57:56 +0300 Subject: [PATCH 3/5] use short timeout in test --- Lib/test/test_external_inspection.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Lib/test/test_external_inspection.py b/Lib/test/test_external_inspection.py index c6332ec311fac1..074dfbaaa90614 100644 --- a/Lib/test/test_external_inspection.py +++ b/Lib/test/test_external_inspection.py @@ -507,7 +507,7 @@ async def main(): def test_all_awaited_by_covers_every_interpreter(self): # gh-158880 async def main_worker(): - await asyncio.sleep(2) + await asyncio.sleep(SHORT_TIMEOUT) async def main(): with InterpreterPoolExecutor() as pool: From 2f8283492d5765d5a0697d8d410bc0370e54918e Mon Sep 17 00:00:00 2001 From: Timofey Ivankov Date: Wed, 7 Oct 2026 11:34:10 +0300 Subject: [PATCH 4/5] rewrite subinterp traversal --- Modules/_remote_debugging/module.c | 27 ++------------------------- 1 file changed, 2 insertions(+), 25 deletions(-) diff --git a/Modules/_remote_debugging/module.c b/Modules/_remote_debugging/module.c index e6056439a11746..a0be653d9bcdc8 100644 --- a/Modules/_remote_debugging/module.c +++ b/Modules/_remote_debugging/module.c @@ -958,11 +958,6 @@ _remote_debugging_RemoteUnwinder_get_all_awaited_by_impl(RemoteUnwinderObject *s if (ensure_async_debug_offsets(self) < 0) { return NULL; } - PyObject *seen = PySet_New(NULL); - if (seen == NULL) { - set_exception_cause(self, PyExc_MemoryError, "Failed to create interpreter set"); - return NULL; - } PyObject *result = PyList_New(0); if (result == NULL) { @@ -971,24 +966,8 @@ _remote_debugging_RemoteUnwinder_get_all_awaited_by_impl(RemoteUnwinderObject *s } // gh-158880: Tasks live in every interpreter, not only the one at the list head - for (uintptr_t interp = self->interpreter_addr; interp != 0; ) { - PyObject *addr = PyLong_FromUnsignedLongLong(interp); - if (addr == NULL) { - set_exception_cause(self, PyExc_MemoryError, "Failed to create interpreter address"); - goto result_err; - } - Py_ssize_t seen_count = PySet_GET_SIZE(seen); - int marked = PySet_Add(seen, addr); - Py_DECREF(addr); - if (marked < 0) { - set_exception_cause(self, PyExc_RuntimeError, "Failed to mark interpreter as seen"); - goto result_err; - } - if (PySet_GET_SIZE(seen) == seen_count) { - // already walked - break; - } - + uintptr_t interp = self->interpreter_addr; + while (interp != 0) { if (refresh_generation_caches_for_interpreter(self, interp) < 0) { goto result_err; } @@ -1023,12 +1002,10 @@ _remote_debugging_RemoteUnwinder_get_all_awaited_by_impl(RemoteUnwinderObject *s } _Py_RemoteDebug_ClearCache(&self->handle); - Py_DECREF(seen); return result; result_err: _Py_RemoteDebug_ClearCache(&self->handle); - Py_DECREF(seen); Py_XDECREF(result); return NULL; } From 6430f21247dd98cd15ab2c4bfb055baadb903dab Mon Sep 17 00:00:00 2001 From: Timofey Ivankov Date: Wed, 7 Oct 2026 11:42:29 +0300 Subject: [PATCH 5/5] rewrite test with queues --- Lib/test/test_external_inspection.py | 34 +++++++++++++++++++--------- 1 file changed, 23 insertions(+), 11 deletions(-) diff --git a/Lib/test/test_external_inspection.py b/Lib/test/test_external_inspection.py index 074dfbaaa90614..299295a18e4c7f 100644 --- a/Lib/test/test_external_inspection.py +++ b/Lib/test/test_external_inspection.py @@ -215,11 +215,17 @@ def _cleanup_sockets(*sockets): ) -def _asyncio_in_subinterpreter(): +def _asyncio_in_subinterpreter(ready, release): + """Park an asyncio task in a subinterpreter until released.""" import asyncio + def wait(loop): + # Signal via the loop, so the task has already suspended + loop.call_soon_threadsafe(ready.put, None) + release.get() + async def sub_worker(): - await asyncio.sleep(2) + await asyncio.to_thread(wait, asyncio.get_running_loop()) asyncio.run(sub_worker()) @@ -506,18 +512,24 @@ async def main(): @requires_subinterpreters def test_all_awaited_by_covers_every_interpreter(self): # gh-158880 + ready = interpreters.create_queue() + release = interpreters.create_queue() + async def main_worker(): await asyncio.sleep(SHORT_TIMEOUT) async def main(): with InterpreterPoolExecutor() as pool: - loop = asyncio.get_running_loop() - loop.run_in_executor(pool, _asyncio_in_subinterpreter) - task = asyncio.create_task(main_worker(), name="main_worker") - self.addCleanup(task.cancel) - for _ in busy_retry(SHORT_TIMEOUT): + try: + loop = asyncio.get_running_loop() + loop.run_in_executor(pool, _asyncio_in_subinterpreter, + ready, release) + task = asyncio.create_task(main_worker(), + name="main_worker") + self.addCleanup(task.cancel) await asyncio.sleep(0) - stacks = [ + ready.get(timeout=SHORT_TIMEOUT) + return [ [frame.funcname.rpartition(".")[2] for frame in coro.call_stack] for info in RemoteUnwinder( @@ -525,11 +537,11 @@ async def main(): for task in info.awaited_by for coro in task.coroutine_stack ] - if ["sleep", "sub_worker"] in stacks: - return stacks + finally: + release.put(None) stacks = asyncio.run(main()) - self.assertIn(["sleep", "sub_worker"], stacks) + self.assertIn(["to_thread", "sub_worker"], stacks) self.assertIn(["sleep", "main_worker"], stacks) @skip_if_not_supported