ashb commented on code in PR #73560:
URL: https://github.com/apache/airflow/pull/73560#discussion_r4104096934


##########
task-sdk/src/airflow/sdk/bases/operator.py:
##########
@@ -200,25 +200,67 @@ 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")
+
+    runner_class = getattr(asyncio, "Runner", None)
+    if runner_class is not None:
+        with runner_class() as runner:
+            yield runner.get_loop()
+        return

Review Comment:
   Nit:
   ```suggestion
       if Runner := getattr(asyncio, "Runner", None):
           with Runner() as runner:
               yield runner.get_loop()
           return
   ```



##########
task-sdk/tests/task_sdk/bases/test_operator.py:
##########
@@ -1137,3 +1140,145 @@ def __init__(self, arg1, arg2, arg3, **kwargs):
     assert op.arg2 == "b"
     assert op.arg3 == 3
     assert op.queue == "THIS"
+
+
+class TestEventLoop:

Review Comment:
   Isn't there an existing test class of BaseAsyncOperator? I generally prefer 
to keep the tests about what class/subject we are testing, not what behaviour 
we are testing.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to