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 0cf78267ad9 Raise DeadlockImminentError for sync comms calls from a
paused event loop thread (#73521)
0cf78267ad9 is described below
commit 0cf78267ad9ff1d4ad73d8eae84668cfb7f381e2
Author: David Blain <[email protected]>
AuthorDate: Fri Sep 25 13:29:14 2026 +0200
Raise DeadlockImminentError for sync comms calls from a paused event loop
thread (#73521)
A synchronous SUPERVISOR_COMMS.send() made from the thread that drives an
event loop, while that loop is paused between two run_until_complete()
calls and an asend() coroutine is parked mid-I/O, blocked forever: the
holder of the comms thread lock can only release it once the loop runs
again, and the loop cannot run while its thread is blocked in send().
The existing guard only covered a *running* loop: it relied on
asyncio.get_running_loop(), which raises while the loop is paused, so
send() fell through to a blocking acquire and the task process froze
with every thread idle, no exception and no log line, until the
execution timeout. This was hit in production by an iterated task whose
executor pulled the next sub-task's input through a sync XCom read
between two run_until_complete() calls.
Treat an in-flight asend() on the current thread's loop as an imminent
deadlock regardless of whether the loop is running or paused, so the
call fails eagerly with a clear DeadlockImminentError instead of hanging.
---
task-sdk/src/airflow/sdk/execution_time/comms.py | 44 +++++++++++++---
.../tests/task_sdk/execution_time/test_comms.py | 60 ++++++++++++++++++++++
2 files changed, 98 insertions(+), 6 deletions(-)
diff --git a/task-sdk/src/airflow/sdk/execution_time/comms.py
b/task-sdk/src/airflow/sdk/execution_time/comms.py
index ee7aa11f331..3ac67112c53 100644
--- a/task-sdk/src/airflow/sdk/execution_time/comms.py
+++ b/task-sdk/src/airflow/sdk/execution_time/comms.py
@@ -122,7 +122,13 @@ ReceiveMsgType = TypeVar("ReceiveMsgType", bound=BaseModel)
class DeadlockImminentError(BaseException):
"""
- Raised when ``send()`` is called from the event loop thread while
``asend()`` holds the lock.
+ Raised when ``send()`` is called from the event loop thread while an
``asend()`` is in flight.
+
+ An in-flight ``asend()`` holds (or is waiting to take) the channel's
thread lock and can only
+ release it once its event loop runs again. That loop runs on the calling
thread, so a blocking
+ ``send()`` there can never be satisfied: the loop cannot run while the
thread is blocked. This
+ holds whether the loop is currently running (a sync SDK call from inside a
coroutine) or paused
+ between two ``run_until_complete`` calls (a sync SDK call from the code
that drives the loop).
Inherits from :class:`BaseException` so it escapes
``contextlib.suppress(Exception)``
and always surfaces to the caller.
@@ -136,9 +142,11 @@ class DeadlockImminentError(BaseException):
def __str__(self) -> str:
return (
f"comms.send() called from the event loop thread for message
'{self.msg_type}' "
- "— deadlock is imminent (asend() is concurrently in-flight). "
+ "— deadlock is imminent (asend() is concurrently in-flight and can
only complete once "
+ "the event loop runs again on this thread). "
"Likely cause: BaseHook.get_hook() or BaseHook.get_connection()
was called "
- "from inside an async task. "
+ "from inside an async task, or a synchronous SDK call was made
from the thread driving "
+ "the event loop while it was paused between two
run_until_complete() calls. "
"Use the async equivalents instead: "
"await BaseHook.aget_hook() or await BaseHook.aget_connection()."
f"\nOffending call stack:\n{self.stack}"
@@ -247,13 +255,37 @@ class CommsDecoder(Generic[ReceiveMsgType, SendMsgType]):
return bool(asyncio.get_running_loop())
return False
+ @property
+ def _asend_in_flight_on_this_thread(self) -> bool:
+ """
+ Whether an ``asend()`` coroutine of a loop that runs on the current
thread is in flight.
+
+ ``asend()`` holds ``_async_lock`` for its whole duration, from before
it takes
+ ``_thread_lock`` (in an executor thread) until after it releases it,
so the async lock
+ being held is the reliable signal that the thread lock is, or is about
to be, held by a
+ coroutine. Such a coroutine only makes progress when its loop runs,
and that loop runs on
+ this thread. This is thread-scoped on purpose: a loop running on
*another* thread keeps
+ going while this one blocks, so ``send()`` from there can safely wait
for the lock.
+ """
+ return threading.get_ident() == self._loop_thread_id and
self._async_lock.locked()
+
def send(self, msg: SendMsgType) -> ReceiveMsgType | None:
"""Send a request to the parent and block until the response is
received."""
frame_bytes = self._make_frame(msg).as_bytes()
- # When called from the event loop thread, use non-blocking acquire to
detect
- # an imminent deadlock: an asend() coroutine currently holds
_thread_lock and
- # is waiting for the event loop to complete its I/O.
+ # An asend() in flight on this thread's loop can only release
_thread_lock once that loop
+ # runs again, and the loop cannot run while this thread blocks in
send(). Raise instead of
+ # blocking. This covers the loop being paused between two
run_until_complete() calls, where
+ # the running-loop check below does not apply:
asyncio.get_running_loop() raises, so send()
+ # would otherwise fall through to a blocking acquire and freeze the
process with every
+ # thread idle (seen with a sync XCom read pulling the next iterated
sub-task's input while a
+ # sub-task's asend() was parked mid-I/O).
+ if self._asend_in_flight_on_this_thread:
+ raise DeadlockImminentError(msg)
+
+ # When called from the event loop thread while the loop is running,
use non-blocking
+ # acquire to detect an imminent deadlock: an asend() coroutine
currently holds
+ # _thread_lock and is waiting for the event loop to complete its I/O.
if not self._thread_lock.acquire(blocking=not self._is_on_loop_thread):
raise DeadlockImminentError(msg)
try:
diff --git a/task-sdk/tests/task_sdk/execution_time/test_comms.py
b/task-sdk/tests/task_sdk/execution_time/test_comms.py
index 665e5418f5b..f328380106a 100644
--- a/task-sdk/tests/task_sdk/execution_time/test_comms.py
+++ b/task-sdk/tests/task_sdk/execution_time/test_comms.py
@@ -17,6 +17,7 @@
from __future__ import annotations
+import asyncio
import threading
import uuid
@@ -341,6 +342,65 @@ class TestCommsDecoder:
assert result is not None
+ def
test_send_from_paused_event_loop_thread_raises_when_asend_in_flight(self,
socket_pair):
+ """
+ Regression: send() called from the loop's own thread while the loop is
*paused* (between
+ two run_until_complete() calls) and an asend() is parked mid-I/O must
raise
+ DeadlockImminentError instead of blocking forever.
+
+ The running-loop check alone does not cover this:
asyncio.get_running_loop() raises while
+ the loop is paused, so send() fell through to a blocking acquire of
_thread_lock. The
+ holder, an asend() coroutine, can only release it once the loop runs
again, on this very
+ thread, which is now blocked. Seen in production as an
IterableOperator freezing with
+ every thread idle: AsyncAwareExecutor.map pulled the next sub-task's
input through a sync
+ XCom read between two run_until_complete() calls while a sub-task's
asend() was in flight.
+ """
+ r, w = socket_pair
+ decoder = CommsDecoder(socket=r, log=structlog.get_logger())
+
+ def _read_request(sock) -> _RequestFrame:
+ length = int.from_bytes(sock.recv(4), "big")
+ body = b""
+ while len(body) < length:
+ body += sock.recv(length - len(body))
+ return msgspec.msgpack.decode(body, type=_RequestFrame)
+
+ def _respond(sock, req: _RequestFrame) -> None:
+ assert req.body is not None
+ resp = {"type": "VariableResult", "key": req.body["key"], "value":
"v"}
+ encoded = msgspec.msgpack.encode(_ResponseFrame(req.id, resp,
None))
+ sock.sendall(len(encoded).to_bytes(4, "big") + encoded)
+
+ loop = asyncio.new_event_loop()
+ try:
+ # Park an asend() mid-I/O: its request gets written, but no
response is sent yet, so
+ # the coroutine sits in the thread reading the response while
holding _thread_lock.
+ in_flight =
loop.create_task(decoder.asend(GetVariable(key="parked")))
+ parked_request =
loop.run_until_complete(asyncio.to_thread(_read_request, w))
+ assert parked_request.body["key"] == "parked"
+ assert not in_flight.done()
+ assert decoder._thread_lock.locked()
+
+ # The loop is now paused and this thread is its thread. Before the
fix this call
+ # blocked forever on _thread_lock.
+ with pytest.raises(DeadlockImminentError) as exc_info:
+ decoder.send(GetVariable(key="should_fail"))
+ assert "deadlock is imminent" in str(exc_info.value)
+
+ # Let the parked asend() finish; the channel must be fully usable
afterwards.
+ _respond(w, parked_request)
+ result = loop.run_until_complete(asyncio.wait_for(in_flight,
timeout=5))
+ assert result.key == "parked"
+
+ server = threading.Thread(target=lambda: _respond(w,
_read_request(w)), daemon=True)
+ server.start()
+ result = decoder.send(GetVariable(key="after"))
+ server.join(timeout=2)
+ assert result.key == "after"
+ finally:
+ loop.close()
+ asyncio.set_event_loop(None)
+
def test_read_frame_recovers_from_short_read_on_header(self):
msg = VariableResult(key="k", value="v", type="VariableResult")
payload = msgspec.msgpack.encode(_ResponseFrame(0, msg.model_dump(),
None))