dabla opened a new pull request, #74074:
URL: https://github.com/apache/airflow/pull/74074

   `dag.test()` (and `dag_maker` in tests) serves a task's Task SDK requests in 
the task's own process through `InProcessSupervisorComms`. That only worked 
while one thread made SDK calls at a time:
   
   - Serving a request removed `task_runner.SUPERVISOR_COMMS` from the whole 
process, so the server-side code (`models.Variable`, `models.Connection`, mask 
forwarding, the secrets backend choice) would not send a request of its own. 
Overlapping calls could leave it removed for good. Every later SDK call in the 
task then failed with `ImportError: cannot import name 'SUPERVISOR_COMMS'`.
   - Answers were matched to requests by order only, so with two threads 
sending at once, one could take the other's answer.
   
   A task that uses a `ThreadPoolExecutor` to fetch Variables or push XComs 
under `dag.test()` hit both.
   
   The comms are no longer removed. A `ContextVar` marks the code serving a 
request: the in-process execution API sets it around each request it serves, so 
its event-loop task and the worker threads of sync routes inherit it, while the 
task's own threads keep their comms. A lock serves one request at a time, from 
handling it to taking its answer.
   
   A per-thread flag was not enough: the in-process API answers on its own 
threads, where a route reading `models.Variable` would send a request of its 
own and wait for the lock forever. With the `ContextVar` the provider tests 
that run tasks under `dag_maker` complete 
(`providers/standard/.../test_python.py`).
   
   Tests:
   - `TestInProcessSupervisorCommsAcrossThreads`: concurrent senders get their 
own answers; only the serving code stops seeing the comms; mask forwarding.
   - `TestInTaskExecutionContext`: the new helper in `airflow.utils.helpers`.
   - `test_app.py`: the in-process app marks the requests it serves.
   - `TestInProcessSupervisorFromThreads`: DB tests running a task whose 
threads push and read XComs and read a DB-stored Variable. They fail on main (5 
out of 5 runs) and pass with this change.
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes (please specify the tool below)
   
   Generated-by: Claude Code (Opus 5.5) following [the 
guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions)
   
   ---
   
   * Read the **[Pull Request 
Guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#pull-request-guidelines)**
 for more information. Note: commit author/co-author name and email in commits 
become permanently public when merged.
   * For fundamental code changes, an Airflow Improvement Proposal 
([AIP](https://cwiki.apache.org/confluence/display/AIRFLOW/Airflow+Improvement+Proposals))
 is needed.
   * When adding dependency, check compliance with the [ASF 3rd Party License 
Policy](https://www.apache.org/legal/resolved.html#category-x).
   * For significant user-facing changes create newsfragment: 
`{pr_number}.significant.rst`, in 
[airflow-core/newsfragments](https://github.com/apache/airflow/tree/main/airflow-core/newsfragments).
 You can add this file in a follow-up commit after the PR is created so you 
know the PR number.
   


-- 
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