Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) by vatsrahul1001 · Pull Request #67204 · apache/airflow · GitHub
Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' [v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) by vatsrahul1001 · Pull Request #67204 · apache/airflow · GitHub
Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' [v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) by vatsrahul1001 · Pull Request #67204 · apache/airflow · GitHub
Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' [v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) by vatsrahul1001 · Pull Request #67204 · apache/airflow · GitHub
Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' [v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) by vatsrahul1001 · Pull Request #67204 · apache/airflow · GitHub
Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' [v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) by vatsrahul1001 · Pull Request #67204 · apache/airflow · GitHub
Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); [v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) by vatsrahul1001 · Pull Request #67204 · apache/airflow · GitHub
Skip to content

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574) - #67204

Merged
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test
May 20, 2026
Merged

[v3-2-test] Recover stuck TIs when direct terminal-state API call fails (#66574)#67204
vatsrahul1001 merged 2 commits into
v3-2-testfrom
backport-173c2a1-v3-2-test

Conversation

@vatsrahul1001

Copy link
Copy Markdown
Contributor

Manual backport of #66574 — auto-backport failed due to the RetryTask branch in supervisor's _handle_request.

Conflict resolution

task-sdk/src/airflow/sdk/execution_time/supervisor.py — one block under elif isinstance(msg, RetryTask):. v3-2-test still had the inline self.client.task_instances.retry(id=..., end_date=..., rendered_map_index=...) call. The PR's squash-merge bundled the "refactor terminal-state dispatch" reshape that replaces all four direct API calls (SucceedTask / RetryTask / DeferTask / RescheduleTask) with self._send_terminal_state_msg(msg). The new _send_terminal_state_msg helper is included in this cherry-pick (defined at supervisor.py:1246 post-merge), so accepting the incoming side at the conflict is the correct resolution — it routes RetryTask through the same dispatcher the other three terminal states already use post-rebase.

The test file (test_supervisor.py) auto-merged cleanly — the new parametrized test_terminal_state_not_set_when_direct_api_fails and test_update_task_state_replays_pending_terminal_state_call (across all four message types) landed as-is.

Scope

2 files changed, +219 / −22 — exact mirror of #66574's merged diff.

* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
…test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
@vatsrahul1001
vatsrahul1001 merged commit 635fc94 into v3-2-testMay 20, 2026
89 checks passed
@vatsrahul1001
vatsrahul1001 deleted the backport-173c2a1-v3-2-test branch May 20, 2026 03:59
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 20, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
vatsrahul1001 added a commit that referenced this pull request May 21, 2026
…ls (#66574) (#67204)
* Recover stuck TIs when direct terminal-state API call fails (#66574)
* Recover stuck TIs when direct terminal-state API call fails
The supervisor's _handle_request for SucceedTask, RetryTask, DeferTask,
and RescheduleTask set _terminal_state BEFORE calling the matching
client.task_instances.{succeed,retry,defer,reschedule}() API. If that
API call raised (transient network blip, server 5xx, etc.),
_terminal_state was set on the supervisor but the server never saw
the transition. The supervisor's update_task_state_if_needed then
saw final_state in STATES_SENT_DIRECTLY and short-circuited the
recovery finish() call -- leaving the TaskInstance stuck RUNNING
on the server forever, blocking downstream dependencies and
triggering false alerts.
Two-part fix:
1. Make the direct API call FIRST. Only set _terminal_state and the
new _terminal_state_synced_to_server flag after the call returns
successfully. If the API raises, both stay unset and the exception
propagates to handle_requests, where the existing catch-all sends
an ErrorResponse to the task subprocess.
2. Have update_task_state_if_needed always call finish() when
_terminal_state_synced_to_server is False, regardless of what
final_state happens to return. The finish() API takes the state
value, so a SUCCESS / DEFERRED / etc. transition that originally
failed is re-attempted via finish() on subprocess exit.
Pre-existing semantics for the no-direct-API states (FAILED,
UP_FOR_RETRY without RetryTask, etc.) preserved -- those land in
the same finish() branch.
Tests added:
- _terminal_state not set when succeed() raises.
- update_task_state_if_needed calls finish() when synced flag is
False, even with final_state == SUCCESS.
- update_task_state_if_needed skips finish() when synced flag is
True (preserves the existing happy-path optimisation).
Reported by the L3 ASVS sweep at apache/tooling-agents#24 (FINDING-007).
* Refactor terminal-state dispatch and parametrize tests across all 4 states
Address review feedback on #66574:
- Extract `_send_terminal_state_msg` helper so the per-msg-type dispatch
for succeed / retry / defer / reschedule lives in one place. Both
`_handle_request` and `_replay_pending_terminal_state_msg` now go
through it instead of duplicating the four-branch isinstance chain.
- Parametrize the two recovery tests over all four terminal-state
message types (was only Succeed + Defer); add UP_FOR_RETRY and
UP_FOR_RESCHEDULE coverage.
* Narrow _pending_terminal_state_msg type to satisfy mypy
The field was annotated as BaseModel | None, but _send_terminal_state_msg
expects SucceedTask | RetryTask | DeferTask | RescheduleTask. mypy
couldn't prove the narrowing at the _replay_pending_terminal_state_msg
call site. Tighten the field type to the exact union the setter assigns
and the consumer accepts.
---------
Co-authored-by: vatsrahul1001 <rah.sharma11@gmail.com>
Co-authored-by: Rahul Vats <43964496+vatsrahul1001@users.noreply.github.com>
(cherry picked from commit 173c2a1)
* Don't pass retry_delay_seconds/retry_reason to retry() — not in v3-2-test signature
The cherry-picked _send_terminal_state_msg dispatcher passed
retry_delay_seconds and retry_reason kwargs (via getattr defensive
fallback) to client.task_instances.retry(). The v3-2-test version of
retry() in task-sdk/src/airflow/sdk/api/client.py only accepts
(id, end_date, rendered_map_index) — those kwargs don't exist on this
branch yet.
In mock-based unit tests the extra kwargs were silently accepted by
Mock but tripped assert_called_once_with. In real DB tests (Postgres
test_task_instance_history_is_created_when_ti_goes_for_retry,
MySQL/SQLite equivalents) the retry() call raised TypeError, the API
server never received the retry transition, and TaskInstanceHistory
never got created — the test's UUID-rotation assertion failed.
---------
Co-authored-by: Jarek Potiuk <jarek@potiuk.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:task-sdktype:bug-fixChangelog: Bug Fixes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@vatsrahul1001@potiuk