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]

Reply via email to