diff --git a/Lib/test/test_external_inspection.py b/Lib/test/test_external_inspection.py index 9fff8e8bff91f7..299295a18e4c7f 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,21 @@ def _cleanup_sockets(*sockets): ) +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.to_thread(wait, asyncio.get_running_loop()) + + 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 +508,42 @@ 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 + ready = interpreters.create_queue() + release = interpreters.create_queue() + + async def main_worker(): + await asyncio.sleep(SHORT_TIMEOUT) + + async def main(): + with InterpreterPoolExecutor() as pool: + 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) + ready.get(timeout=SHORT_TIMEOUT) + 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 + ] + finally: + release.put(None) + + stacks = asyncio.run(main()) + self.assertIn(["to_thread", "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..a0be653d9bcdc8 100644 --- a/Modules/_remote_debugging/module.c +++ b/Modules/_remote_debugging/module.c @@ -958,9 +958,6 @@ _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) { - return NULL; - } PyObject *result = PyList_New(0); if (result == NULL) { @@ -968,23 +965,40 @@ _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 + uintptr_t interp = self->interpreter_addr; + while (interp != 0) { + if (refresh_generation_caches_for_interpreter(self, interp) < 0) { + goto result_err; + } - uintptr_t head_addr = self->interpreter_addr - + (uintptr_t)self->async_debug_offsets.asyncio_interpreter_state.asyncio_tasks_head; + // Process all threads + if (iterate_threads(self, interp, process_thread_for_awaited_by, result) < 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; + 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); @@ -1064,7 +1078,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)) {