dabla opened a new pull request, #73521: URL: https://github.com/apache/airflow/pull/73521
### Problem While running an iterated task in production (Airflow 3.3.2), the task process froze a few times a week: a batch of sub-tasks logged "Attempting running task" and nothing else, every thread showed zero context switches, and the supervisor kept heartbeating until the 24h execution timeout. No exception, no log line, nothing on the API server. A `faulthandler` dump of a frozen process showed the deadlock: 1. The code driving the event loop pulled the next sub-task's input on the main thread **between two `run_until_complete()` calls**, i.e. with the loop paused. 2. That pull was a synchronous supervisor call (`LazyXComSequence.__getitem__` -> `SUPERVISOR_COMMS.send()`). With no loop running, `send()` acquired the comms thread lock in **blocking** mode. 3. A sub-task coroutine parked mid-`asend()` already held that lock and could only release it once the loop ran again. The loop runs on the very thread now blocked in `send()`, so it never did. The SDK already guards against a sync `send()` from a *running* loop (#68377 raises `DeadlockImminentError`), but that guard relies on `asyncio.get_running_loop()`, which raises while the loop is paused. A sync call from the loop's thread while the loop is paused therefore fell through to a blocking acquire instead of raising, and anything in flight on the loop was stuck behind it. ### Fix `CommsDecoder.send()` now treats an in-flight `asend()` on the current thread's loop as an imminent deadlock whether the loop is running or paused, and raises `DeadlockImminentError` eagerly instead of blocking. The signal is `_async_lock.locked()`: `asend()` holds the async lock for its whole duration, from before it takes the thread lock (in an executor thread) until after it releases it, so it is the reliable indicator that the thread lock is, or is about to be, held by a coroutine that needs this thread's loop to make progress. The check stays 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. The existing running-loop check is kept as is; the error message now mentions the paused-loop cause. Every previously-hanging call now fails with a clear error naming the message type and the offending call stack, while sync sends from the loop thread with no `asend()` in flight (e.g. the task runner pushing XComs after `execute()` returns) keep working as before, which the existing tests cover. ### Testing * New regression test `test_send_from_paused_event_loop_thread_raises_when_asend_in_flight` parks an `asend()` mid-I/O, pauses the loop, and asserts a sync `send()` from the loop's thread raises. Verified it hangs (killed by a 20s timeout) without the fix and passes with it. It also checks the channel is fully usable once the parked `asend()` completes. * `tests/task_sdk/execution_time/test_comms.py` and `test_context.py`: 203 passed. * `prek` pre-commit stage on the changed files, including `mypy-task-sdk`, `ruff` and `check-supervisor-schemas-versions`: passed. --- ##### Was generative AI tooling used to co-author this PR? - [X] Yes — Claude Code (Fable 5.1) Generated-by: Claude Code (Fable 5.1) following [the guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions) -- 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]
