Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changelog/70226.fixed.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Stopped proxy minions logging ``Error during asyncio shutdown: The future belongs to a different loop than the one specified as the loop argument`` on Python 3.14. ``asyncio.gather`` takes no ``loop`` argument any more and resolves the loop from the calling context, but ``SyncWrapper.close()`` runs outside the loop it is tearing down, so gathering that loop's pending tasks was rejected and they were never drained.
29 changes: 21 additions & 8 deletions salt/utils/asynchronous.py
Original file line number Diff line number Diff line change
Expand Up @@ -184,16 +184,29 @@ def close(self):
if pending_tasks:
for task in pending_tasks:
task.cancel()
gathered = asyncio.gather(*pending_tasks, return_exceptions=True)

# ``asyncio.gather`` has no ``loop`` argument any more, so it
# resolves the loop from the calling context. ``close()``
# runs outside ``self.asyncio_loop`` -- the thread's current
# loop is a different one -- so on Python 3.14 gathering
# tasks that belong to ``self.asyncio_loop`` raises
# ``ValueError: The future belongs to a different loop than
# the one specified as the loop argument``. Earlier versions
# took the loop from the first future and let it pass.
#
# Build the gather *inside* the loop instead, where the
# running loop is the right one on every version.
async def _drain(tasks):
await asyncio.gather(*tasks, return_exceptions=True)

drain = _drain(pending_tasks)
try:
self.asyncio_loop.run_until_complete(gathered)
self.asyncio_loop.run_until_complete(drain)
except Exception: # pylint: disable=broad-except
# ``gathered`` is a Future; if run_until_complete bailed
# part-way we still need to make sure the Future is
# consumed so its exception (if any) isn't logged as
# unhandled. Tasks already cancelled above.
if not gathered.done():
gathered.cancel()
# Close the coroutine we just built so it is not
# garbage-collected unawaited, which would emit a
# RuntimeWarning on stderr. Tasks already cancelled.
drain.close()

if self._loop_can_run_until_complete(self.asyncio_loop):
shutdown_agens = self.asyncio_loop.shutdown_asyncgens()
Expand Down
63 changes: 63 additions & 0 deletions tests/pytests/unit/utils/test_asynchronous.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
import tornado.ioloop

import salt.utils.asynchronous as asynchronous
from tests.support.mock import patch


class HelperA:
Expand Down Expand Up @@ -156,3 +157,65 @@ def test_sync_wrapper_thread_has_asyncio_loop_65702():
assert sync.check_loop() is True
finally:
sync.close()


class HelperPending:
"""A helper whose wrapped coroutine leaves a task pending on the loop."""

async_methods = [
"start_background",
]

def __init__(self, io_loop=None):
self.io_loop = io_loop

@tornado.gen.coroutine
def start_background(self):
# Leave a long-lived task behind on this wrapper's own loop, so
# ``close()`` has something to drain.
asyncio.ensure_future(asyncio.sleep(3600))
raise tornado.gen.Return(True)


def test_close_drains_tasks_belonging_to_the_wrappers_own_loop():
"""
``close()`` runs outside the loop it is tearing down -- the calling
thread's current loop is a different one. ``asyncio.gather`` no longer
takes a ``loop`` argument, so it resolves the loop from the calling
context, and on Python 3.14 gathering tasks that belong to another loop
raises ``ValueError: The future belongs to a different loop than the one
specified as the loop argument``. Earlier versions took the loop from the
first future and let it through, so this surfaced as a wall of
"Error during asyncio shutdown" for every proxy minion on 3.14.

Building the gather inside the loop drains the tasks on every version.
"""
sync = asynchronous.SyncWrapper(HelperPending)
sync.start_background()

pending = [t for t in asyncio.all_tasks(sync.asyncio_loop) if not t.done()]
assert pending, "expected a task pending on the wrapper's loop"

# The failure only happens when ``close()`` is called from inside a
# *different running* loop, which is how it is reached in a proxy minion:
# ``asyncio.gather`` then resolves the running loop rather than the tasks'
# own loop and rejects them. Drive it that way.
#
# Asserting on the tasks alone would not catch this either -- they are
# cancelled before the gather, so they end up done() regardless. The
# symptom is the swallowed exception, so assert nothing was logged.
async def _close_from_another_running_loop():
with patch.object(asynchronous.log, "error") as log_error:
sync.close()
return log_error.call_args_list

driver = asyncio.new_event_loop()
try:
errors = driver.run_until_complete(_close_from_another_running_loop())
finally:
driver.close()

# Only the swallowed exception is asserted on. The tasks themselves
# cannot be driven to completion here -- a loop cannot be run from inside
# another running loop -- so their state is not the thing under test.
assert not errors, errors
Loading