This is an automated email from the ASF dual-hosted git repository.

ashb pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 8d815efdbc3 Own the BaseAsyncOperator event loop through 
asyncio.Runner (#73560)
8d815efdbc3 is described below

commit 8d815efdbc3dad1ad5cdf592934d1b4de7458164
Author: David Blain <[email protected]>
AuthorDate: Tue Sep 29 15:49:20 2026 +0200

    Own the BaseAsyncOperator event loop through asyncio.Runner (#73560)
    
    Every task that runs through BaseAsyncOperator.execute logs 
DeprecationWarning: There is no current event loop as its first line on Python 
3.12 and 3.13:
    
    ```
    [2026-09-22 16:34:41] INFO - Worker startup parse complete 
bundle_name=dags-folder bundle_version= dag_file=test_msgraph.py 
bundle_prepare_ms=1 dag_file_parse_ms=17
    [2026-09-22 16:34:41] WARNING - There is no current event loop 
category=DeprecationWarning
    ```
    
    The task-runner process starts without an event loop set, and the 
event_loop() helper relied on asyncio.get_event_loop() to create one. That is 
exactly the case Python deprecated in 3.12, and Python 3.14 raises RuntimeError 
there instead of creating a loop. The helper also never closed a loop that 
asyncio created implicitly on its behalf, so that loop leaked for the rest of 
the process.
    
    event_loop() now owns its loop through asyncio.Runner, as the asyncio docs 
recommend in place of policies. Plain asyncio.run() does not cover the helper's 
purpose: it creates a loop for one coroutine and closes it, while the helper 
hands out a loop that outlives a single coroutine so callers can drive it 
through several run_until_complete() calls (the upcoming iterable operator does 
exactly that between batches of sub-tasks). Runner gives that, and on exit also 
cancels leftover tasks  [...]
---
 task-sdk/src/airflow/sdk/bases/operator.py     |  69 ++++++++--
 task-sdk/tests/task_sdk/bases/test_operator.py | 171 +++++++++++++++++++++++++
 2 files changed, 226 insertions(+), 14 deletions(-)

diff --git a/task-sdk/src/airflow/sdk/bases/operator.py 
b/task-sdk/src/airflow/sdk/bases/operator.py
index 940fee84b87..7634443b119 100644
--- a/task-sdk/src/airflow/sdk/bases/operator.py
+++ b/task-sdk/src/airflow/sdk/bases/operator.py
@@ -200,25 +200,66 @@ def coerce_resources(resources: dict[str, Any] | None) -> 
Resources | None:
     return Resources(**resources)
 
 
+def _cancel_all_tasks(loop: AbstractEventLoop) -> None:
+    """Cancel every task still pending on ``loop`` and let it run its cleanup, 
as ``asyncio.run()`` does."""
+    to_cancel = asyncio.all_tasks(loop)
+    if not to_cancel:
+        return
+    for task in to_cancel:
+        task.cancel()
+    loop.run_until_complete(asyncio.gather(*to_cancel, return_exceptions=True))
+    for task in to_cancel:
+        if not task.cancelled() and task.exception() is not None:
+            loop.call_exception_handler(
+                {
+                    "message": "unhandled exception during event_loop() 
shutdown",
+                    "exception": task.exception(),
+                    "task": task,
+                }
+            )
+
+
 @contextlib.contextmanager
 def event_loop() -> Generator[AbstractEventLoop]:
-    new_event_loop = False
-    loop = None
+    """
+    Own an event loop for the duration of the block, to drive coroutines from 
synchronous code.
+
+    Unlike ``asyncio.run()``, which creates a loop for one coroutine and 
closes it, the loop yielded here
+    outlives any number of ``run_until_complete()`` calls: tasks created on it 
in one call are still there
+    for the next, which is what code that pauses and resumes a loop between 
batches of work needs.
+
+    On Python 3.11+ this is :class:`asyncio.Runner`; Python 3.10 has no 
``Runner`` and gets a fallback that
+    replicates ``Runner.close()``. In both cases the loop is set as the 
current loop of the thread while the
+    block runs and unset afterwards, leftover tasks are cancelled, async 
generators and the default executor
+    are shut down and the loop is closed. Nothing goes through 
``asyncio.get_event_loop()``, which since
+    Python 3.12 warns and since 3.14 raises when no loop is set. Entering the 
block from a running loop is
+    refused: a loop cannot drive another loop on the same thread.
+    """
+    try:
+        asyncio.get_running_loop()
+    except RuntimeError:
+        pass
+    else:
+        raise RuntimeError("event_loop() cannot be used from a running event 
loop")
+
+    if Runner := getattr(asyncio, "Runner", None):
+        with Runner() as runner:
+            yield runner.get_loop()
+        return
+
+    # Python 3.10 has no asyncio.Runner: replicate what Runner.close() does. 
Drop with Python 3.10 support.
+    loop = asyncio.new_event_loop()
+    asyncio.set_event_loop(loop)
     try:
-        try:
-            loop = asyncio.get_event_loop()
-            if loop.is_closed():
-                raise RuntimeError
-        except RuntimeError:
-            loop = asyncio.new_event_loop()
-            asyncio.set_event_loop(loop)
-            new_event_loop = True
         yield loop
     finally:
-        if new_event_loop and loop is not None:
-            with contextlib.suppress(AttributeError):
-                loop.close()
-                asyncio.set_event_loop(None)
+        try:
+            _cancel_all_tasks(loop)
+            loop.run_until_complete(loop.shutdown_asyncgens())
+            loop.run_until_complete(loop.shutdown_default_executor())
+        finally:
+            asyncio.set_event_loop(None)
+            loop.close()
 
 
 class _PartialDescriptor:
diff --git a/task-sdk/tests/task_sdk/bases/test_operator.py 
b/task-sdk/tests/task_sdk/bases/test_operator.py
index dcb5240a83d..9041e404b12 100644
--- a/task-sdk/tests/task_sdk/bases/test_operator.py
+++ b/task-sdk/tests/task_sdk/bases/test_operator.py
@@ -17,6 +17,7 @@
 
 from __future__ import annotations
 
+import asyncio
 import copy
 import logging
 import uuid
@@ -32,12 +33,14 @@ import structlog
 from airflow.sdk import DAG, Label, TaskGroup, task as task_decorator
 from airflow.sdk._shared.secrets_masker import _secrets_masker, mask_secret
 from airflow.sdk.bases.operator import (
+    BaseAsyncOperator,
     BaseOperator,
     BaseOperatorMeta,
     ExecutorSafeguard,
     chain,
     chain_linear,
     cross_downstream,
+    event_loop,
 )
 from airflow.sdk.definitions.param import ParamsDict
 from airflow.sdk.definitions.template import literal
@@ -1137,3 +1140,171 @@ def test_partial_default_args():
     assert op.arg2 == "b"
     assert op.arg3 == 3
     assert op.queue == "THIS"
+
+
[email protected](params=["asyncio.Runner", "fallback"])
+def event_loop_runner(request, monkeypatch):
+    """
+    Run a test against both ways ``event_loop()`` can own its loop.
+
+    ``asyncio.Runner`` exists from Python 3.11; the fallback replicates it for 
Python 3.10 and is
+    exercised on every Python by hiding ``asyncio.Runner`` for the test.
+    """
+    if request.param == "fallback":
+        monkeypatch.delattr(asyncio, "Runner", raising=False)
+    elif not hasattr(asyncio, "Runner"):
+        pytest.skip("asyncio.Runner needs Python 3.11+")
+    return request.param
+
+
[email protected]
+def fresh_process_loop_state():
+    """
+    Mimic a freshly started task-runner process: no loop set and 
``set_event_loop()`` never called.
+
+    In that state ``asyncio.get_event_loop()`` emits ``DeprecationWarning: 
There is no current event
+    loop`` when it has to create a loop (Python 3.12+) or raises (Python 
3.14+).
+    """
+    with warnings.catch_warnings():
+        # The policy API is deprecated on Python 3.14; it is still the only 
way to reset this state.
+        warnings.simplefilter("ignore", DeprecationWarning)
+        previous = asyncio.get_event_loop_policy()
+        asyncio.set_event_loop_policy(asyncio.DefaultEventLoopPolicy())
+        try:
+            yield
+        finally:
+            asyncio.set_event_loop_policy(previous)
+
+
+def _deprecation_warnings(caught: list[warnings.WarningMessage]) -> list[str]:
+    return [str(w.message) for w in caught if issubclass(w.category, 
DeprecationWarning)]
+
+
[email protected]("event_loop_runner")
+class TestBaseAsyncOperator:
+    class AsyncOperator(BaseAsyncOperator):
+        async def aexecute(self, context):
+            await asyncio.sleep(0)
+            return "done"
+
+    def test_is_async(self):
+        assert self.AsyncOperator(task_id="async_op").is_async is True
+
+    def test_execute_requires_aexecute(self):
+        with pytest.raises(NotImplementedError):
+            BaseAsyncOperator(task_id="async_op").execute({})
+
+    @pytest.mark.usefixtures("fresh_process_loop_state")
+    @pytest.mark.parametrize("execution_timeout", [None, 
timedelta(seconds=5)], ids=["no-timeout", "timeout"])
+    def test_execute_runs_aexecute_without_deprecation_warning(self, 
execution_timeout):
+        op = self.AsyncOperator(task_id="async_op", 
execution_timeout=execution_timeout)
+
+        with warnings.catch_warnings(record=True) as caught:
+            warnings.simplefilter("always")
+            assert op.execute({}) == "done"
+
+        assert _deprecation_warnings(caught) == []
+        # The loop execute() ran on is closed and not left behind as the 
thread's current loop.
+        with pytest.raises(RuntimeError):
+            asyncio.get_event_loop()
+
+    def test_execute_enforces_execution_timeout(self):
+        class SlowOperator(BaseAsyncOperator):
+            async def aexecute(self, context):
+                await asyncio.sleep(60)
+
+        op = SlowOperator(task_id="slow_op", 
execution_timeout=timedelta(milliseconds=50))
+
+        with pytest.raises(asyncio.TimeoutError):
+            op.execute({})
+
+
[email protected]("event_loop_runner")
+class TestEventLoop:
+    """``event_loop()`` hands synchronous code a loop it owns, through 
``asyncio.Runner`` where Python has it."""
+
+    @pytest.mark.usefixtures("fresh_process_loop_state")
+    def test_creates_and_closes_own_loop_without_deprecation_warning(self):
+        with warnings.catch_warnings(record=True) as caught:
+            warnings.simplefilter("always")
+            with event_loop() as loop:
+                assert not loop.is_closed()
+                assert loop.run_until_complete(asyncio.sleep(0, result="ran")) 
== "ran"
+
+        assert _deprecation_warnings(caught) == []
+        assert loop.is_closed()
+        # Nothing is left behind as the thread's current loop.
+        with pytest.raises(RuntimeError):
+            asyncio.get_event_loop()
+
+    @pytest.mark.usefixtures("fresh_process_loop_state")
+    def test_second_use_after_cleanup_still_works(self):
+        with event_loop() as first:
+            pass
+        with warnings.catch_warnings(record=True) as caught:
+            warnings.simplefilter("always")
+            with event_loop() as second:
+                assert second.run_until_complete(asyncio.sleep(0, 
result="ran")) == "ran"
+
+        assert _deprecation_warnings(caught) == []
+        assert first.is_closed()
+        assert second.is_closed()
+
+    def test_loop_outlives_run_until_complete_calls(self):
+        """Unlike ``asyncio.run()``, work scheduled in one run is still there 
for the next."""
+        with event_loop() as loop:
+            release = asyncio.Event()
+
+            async def wait_for_release():
+                await release.wait()
+                return "released"
+
+            pending = loop.create_task(wait_for_release())
+            assert loop.run_until_complete(asyncio.sleep(0, result="first 
run")) == "first run"
+            assert not pending.done()
+
+            release.set()
+            assert loop.run_until_complete(pending) == "released"
+
+    def test_pending_tasks_are_cancelled_on_exit(self):
+        cleaned_up = False
+
+        async def never_finishes():
+            nonlocal cleaned_up
+            try:
+                await asyncio.sleep(3600)
+            except asyncio.CancelledError:
+                cleaned_up = True
+                raise
+
+        with event_loop() as loop:
+            pending = loop.create_task(never_finishes())
+            # Let the task start so cancellation reaches its body.
+            loop.run_until_complete(asyncio.sleep(0))
+
+        assert pending.cancelled()
+        assert cleaned_up
+        assert loop.is_closed()
+
+    def test_leaves_loop_already_set_open_but_does_not_reuse_it(self):
+        existing = asyncio.new_event_loop()
+        asyncio.set_event_loop(existing)
+        try:
+            with event_loop() as loop:
+                assert loop is not existing
+                assert loop.run_until_complete(asyncio.sleep(0, result="ran")) 
== "ran"
+            assert loop.is_closed()
+            assert not existing.is_closed()
+        finally:
+            existing.close()
+            asyncio.set_event_loop(None)
+
+    @pytest.mark.asyncio
+    async def test_refuses_running_loop(self):
+        running = asyncio.get_running_loop()
+
+        with pytest.raises(RuntimeError, match="running event loop"):
+            with event_loop():
+                pass
+
+        assert not running.is_closed()

Reply via email to