Repository navigation
Test get_async_extra_dejson against the real supervisor comms - #74166
Merged
dabla merged 4 commits intoOct 4, 2026
Merged
Conversation
… 3.0 and 3.1 get_async_extra_dejson() falls back to running extra_dejson in a worker thread when Connection.aextra_dejson() is missing. Its secret masking sends to the supervisor, and on Airflow 3.0 and 3.1 the supervisor channel has no thread lock: a send from a worker thread could interleave with an asend() in flight on the event loop, e.g. in the triggerer. On those versions the helper now deserializes Connection.extra without the supervisor, as async hooks did before (apache#55179, apache#72130). The worker thread stays for Airflow 3.2 to 3.3.1, whose channel is thread-safe, and Airflow 2, which has no supervisor. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 task done
shahar1
approved these changes
Oct 3, 2026
Contributor
|
Please note the failing tests |
Revert the Airflow 3.0/3.1 json.loads branch: on those versions async code only runs in the triggerer, whose channel serializes every request, a worker thread's synchronous send included, through asend() and an asyncio.Lock, so the worker-thread fallback was safe and kept the masking. Instead, add a test that forks a real task process talking to the real supervisor of the installed Airflow and reads extras concurrently from an event loop, so each compat job checks the helper against that version's comms rather than mocks. It needs Airflow 3.2+, the first version where a task process can call the supervisor from async code. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Follow-up to #74147: a test of
get_async_extra_dejson()against the real supervisor comms of the installed Airflow, instead of mocks.The helper exists because the Task SDK's supervisor channel differs between Airflow versions (
asend(), thread and async locks,Connection.aextra_dejson()). A mocked channel only encodes assumptions about those differences; this test forks a real task process talking to the real supervisor (as the Task SDK's own supervisor tests do, only the Execution API client is stubbed) and, from an event loop, resolves a connection and reads its extra in 50 concurrent coroutines over 3 rounds. It asserts the task process succeeds, every read returns the extra, and the decoded extra was masked through the supervisor. As provider tests run against each supported Airflow version, every compat job checks the helper against that version's comms.It runs on Airflow 3.2+, the first version where a task process can call the supervisor from async code: before that, task-process
asend()is not implemented (3.1) or absent (3.0), and async code only runs in the triggerer, whose channel serializes every request — a worker thread's synchronous send included — throughasend()and anasyncio.Lock.Verified locally:
extra_dejsonsynchronously on the loopDeadlockImminentError(MaskSecret)An earlier version of this PR changed the helper to skip the worker-thread fallback on Airflow 3.0/3.1; running against the real comms showed that fallback is safe there (see above), so that change was reverted and only the test remains. The forked child also resets asgiref's shared
sync_to_asyncexecutor, which it would otherwise inherit from the pytest process (started by earlier tests) without its thread.Was generative AI tooling used to co-author this PR?
Generated-by: Claude Opus 5.5 following the guidelines