Skip to content

Run loop gates and advance iterations atomically - #74348

Open
ashb wants to merge 2 commits into
task-loops-stack-6from
task-loops-stack-7
Open

ashb wants to merge 2 commits into
task-loops-stack-6from
task-loops-stack-7

Conversation

@ashb

@ashb ashb commented Oct 6, 2026

Copy link
Copy Markdown
Member

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

  • Read the 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) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

Comment thread airflow-core/src/airflow/models/task_coordinates.py Outdated
Comment thread airflow-core/src/airflow/models/dagrun.py Outdated
Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py Outdated
Comment thread airflow-core/src/airflow/jobs/scheduler_job_runner.py
Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py
Comment thread airflow-core/src/airflow/models/task_coordinates.py Outdated
Comment thread airflow-core/tests/unit/models/test_task_coordinates.py Outdated
@ashb
ashb force-pushed the task-loops-stack-7 branch from 43debd4 to 3f2fff3 Compare October 7, 2026 13:57
@ashb
ashb removed this pull request from stack #74339 October 7, 2026 15:20
@ashb
ashb force-pushed the task-loops-stack-7 branch from 3f2fff3 to 2cd460c Compare October 7, 2026 15:22
@ashb
ashb added this pull request to stack #74410 October 7, 2026 15:22
@ashb
ashb force-pushed the task-loops-stack-7 branch 2 times, most recently from 15a1960 to a6d33c8 Compare October 7, 2026 20:57
@ashb
ashb force-pushed the task-loops-stack-7 branch from a6d33c8 to cc4f2a3 Compare October 8, 2026 13:56
@ashb
ashb force-pushed the task-loops-stack-7 branch from cc4f2a3 to 08d8d0c Compare October 8, 2026 16:06
Comment thread devel-common/src/tests_common/test_utils/asserts.py Outdated
Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py
Comment thread airflow-core/src/airflow/models/dagrun.py Outdated
Comment thread task-sdk/tests/task_sdk/execution_time/test_loop.py Outdated
Comment thread airflow-core/src/airflow/models/dynamic_region.py Outdated
Comment thread airflow-core/src/airflow/models/task_coordinates.py Outdated
Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py Outdated
Comment thread airflow-core/src/airflow/models/dagrun.py Outdated
Comment thread task-sdk/tests/task_sdk/execution_time/test_loop.py
@ashb
ashb force-pushed the task-loops-stack-7 branch from 08d8d0c to 636464c Compare October 8, 2026 20:45
@ashb
ashb force-pushed the task-loops-stack-7 branch from 636464c to 240027d Compare October 9, 2026 10:32
@ashb
ashb marked this pull request as ready for review October 9, 2026 12:39
@ashb
ashb requested a review from amoghrajesh as a code owner October 9, 2026 12:39
@ashb
ashb force-pushed the task-loops-stack-7 branch from 240027d to cde2024 Compare October 9, 2026 14:03
@ashb
ashb force-pushed the task-loops-stack-7 branch from cde2024 to 262afa3 Compare October 9, 2026 16:28
@ashb
ashb force-pushed the task-loops-stack-7 branch from 262afa3 to 81a80be Compare October 9, 2026 20:34
@ashb
ashb force-pushed the task-loops-stack-7 branch 2 times, most recently from 8f9fe24 to 8f2f3b2 Compare October 9, 2026 22:18
ashb added 2 commits October 10, 2026 08:04
The scheduler, the API and the UI must never see every task of a loop finished
while its next pass does not exist yet. If we recorded a continuing gate
as successful and created the next pass afterwards, a crash between the two
would leave recovery unable to tell an unfinished advancement from a completed
decision. Gate completion and next-pass creation therefore happen in one
transaction under the DagRun lock, which also stops a duplicate completion from
advancing the loop twice.

The continue-or-stop decision is a reserved XCom stored under the gate's attempt
UUID. A retried gate starts without a decision, and a retired attempt cannot
advance the loop. An invalid decision moves the gate attempt to retry or failed
within the same request; rejecting the request instead would leave the attempt
running with nobody to finish it.

Gates use normal task dependencies and an explicit loop context. A barrier
that waited for every body task would override trigger rules that allow early
progress. Workers and dag.test supply the same coordinates, so decisions and
previous-iteration reads reach the intended gate and data.

Skipping downstream of a task from outside a loop skips every live pass.
Skipping only one would let later passes run work the author meant to skip.
On an unversioned bundle, renaming a loop's gate mid-run re-pins the running
gate attempt to a Dag version that no longer has its task. Its completion then
looked up the loop through that task and the request failed with a 500, leaving
the attempt running with nobody to finish it. Supporting the rename is out of
scope; the attempt is failed instead, without retry, because a retry would hit
the same missing task.
@ashb
ashb force-pushed the task-loops-stack-7 branch from 8f2f3b2 to c834724 Compare October 10, 2026 07:05

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants