Skip to content

Task SDK: don't let exit code override a confirmed terminal state - #73142

Open
seanmuth wants to merge 1 commit into
apache:mainfrom
seanmuth:seanmuth/fix-ti-terminal-state-precedence
Open

seanmuth wants to merge 1 commit into
apache:mainfrom
seanmuth:seanmuth/fix-ti-terminal-state-precedence

Conversation

@seanmuth

@seanmuth seanmuth commented Sep 14, 2026 •

Copy link
Copy Markdown
Contributor

A subprocess can report a terminal state via message (SucceedTask, TaskState(SKIPPED),
etc.) and then exit with a genuinely non-zero code afterward -- e.g. an OOM kill, or
_handle_process_overtime_if_needed() sending SIGTERM once _terminal_state is already
set. Supervisor.final_state previously derived its result from the exit code whenever
one was observed, ignoring an already-confirmed terminal state in that case.

The clearest reachable case is SUCCESS: SucceedTask writes the row directly
(_send_terminal_state_msg) and clears the pending-message slot on success, so a later
non-zero exit code re-triggers update_task_state_if_needed() -> .finish() against a
row succeed() already wrote, which 409s against the server and crashes task
supervision. The same precedence gap is also reachable for TaskState(FAILED) with
retries enabled (should_retry=True): the worker legitimately reports FAILED, but a
later non-zero exit resolves to UP_FOR_RETRY on the old code, so nothing gets written
and the row is left stuck RUNNING instead of FAILED.

Fix: trust self._terminal_state whenever it's already set, and only fall back to
deriving the state from the exit code when no terminal state was ever set at all. This
is a reordering of final_state's existing branches, purely local, no additional
network/DB round-trip.

Rebased over #73249 / #73253's rework of pending-message delivery. Those changes
mean TaskState-reported outcomes (SKIPPED, FAILED, etc.) are now normally dispatched
in update_task_state_if_needed() using the message's own state directly, without
consulting final_state at all while a message is pending -- so for that path, this
precedence is now defense-in-depth rather than the reachable bug it was before that
rework. The SUCCESS case above remains fully live and reachable regardless. This also
subsumes final_state's existing special-case handling of SERVER_TERMINATED and
UP_FOR_RETRY (added by #73253) under one general rule: every terminal state gets the
same precedence, not just those two.

Root-caused via a live instrumented burst-fanout repro that reproduced the SUCCESS
case end to end (full write-up: #65708 (comment)).
This addresses the underlying ordering bug rather than only swallowing its symptom (the
separate, already-merged #63355 idempotency check on the API side remains a good
complementary hardening for the retried-succeed() case it covers).

closes: #65708


Was generative AI tooling used to co-author this PR?
  • Yes — Claude Sonnet 5

Generated-by: Claude Sonnet 5 following the guidelines

Comment thread task-sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment thread task-sdk/src/airflow/sdk/execution_time/supervisor.py
Comment thread task-sdk/tests/task_sdk/execution_time/test_supervisor.py Outdated
Comment thread airflow-core/newsfragments/73142.bugfix.rst Outdated
@seanmuth

seanmuth commented Sep 15, 2026 •

Copy link
Copy Markdown
Contributor Author

Correction to my own earlier framing: the "an unobserved/defaulted exit code overriding an already-confirmed terminal state" explanation in the current docstring is inference I hadn't actually verified. I built live instrumentation and reproduced the crash under load — full writeup with tracebacks and the apiserver-side timeline is in #65708 (comment).

Short version: the real, reachable case is update_task_state_if_needed() computing final_state=FAILED from a genuinely non-zero exit_code despite _terminal_state already being SUCCESS, then calling .finish(), which 409s against the row a prior .succeed() call already wrote correctly — and that exception is uncaught in wait()/supervise(). Traced end to end against this PR's fix: with correct precedence, final_state resolves to SUCCESS, lands in STATES_SENT_DIRECTLY, and .finish() is never called. This PR's change holds up against the live-captured case; the docstring's cited scenario (_handle_process_overtime_if_needed()/SIGTERM) was just the wrong example, not a wrong fix.

One addition since I first posted this: the other duplicate write in that same trace (a retried .succeed() call landing on an already-success row, separate from the .finish() case this PR fixes) turns out to already be fixed too, by PR #63355's idempotency short-circuit in the execution API — also 3.3.0+, also not backported to the branch this was reproduced on. Details in the linked comment.

Holding off on any further changes here pending review.


Drafted-by: Claude Code (Sonnet 5); reviewed by @seanmuth before posting

Comment thread task-sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment thread task-sdk/src/airflow/sdk/execution_time/supervisor.py Outdated
Comment thread task-sdk/tests/task_sdk/execution_time/test_supervisor.py
Comment thread task-sdk/tests/task_sdk/execution_time/test_supervisor.py
@seanmuth
seanmuth force-pushed the seanmuth/fix-ti-terminal-state-precedence branch from 10b490d to 86a3f75 Compare September 17, 2026 17:16
@seanmuth seanmuth changed the title Task SDK: trust a confirmed terminal state over an unobserved exit code Task SDK: don't let exit code override a confirmed terminal state Sep 17, 2026
@seanmuth

Copy link
Copy Markdown
Contributor Author

Pushed a squashed rewrite addressing this round of review:

  • Retitled and reworded the PR body and the single remaining commit to drop the retracted "unobserved exit code defaults to 1" framing entirely, so it no longer survives into squash-merge history. The docstring now cites the two verified reachable cases instead: SUCCESS (row already written by succeed(), a later non-zero exit re-triggers a 409ing .finish()) and SKIPPED (no direct API call, row still RUNNING, so a later non-zero exit either overwrites it with FAILED or -- with retries -- abandons the write entirely).
  • Fixed the test ordering bug: update_task_state_if_needed() and finish.assert_not_called() now run before the final_state assertion, so should_retry=False actually pins the spurious .finish() call instead of pytest stopping at the first (already-informative) assertion.
  • Added test_confirmed_skipped_state_persists_over_later_nonzero_exit_code, parametrized over both should_retry values, asserting finish() is called with state=SKIPPED -- covering the half of the behavior change (states outside STATES_SENT_DIRECTLY) the existing tests never exercised.

Full suite green (211 passed, 1 skipped/platform-only) plus mypy clean after each change.


Drafted-by: Claude Sonnet 5 (no human review before posting)

@kaxil kaxil left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-reviewed at 86a3f75. Everything from the earlier rounds is addressed and the precedence change is correct in every (_terminal_state, _exit_code, _should_retry) combination I traced. Approving.

Three things to take or leave, none worth another round:

supervisor.py:1813-1816 states the 409 unconditionally, but it only holds for should_retry=False. With retries the pre-fix code resolved to UP_FOR_RETRY, which is in STATES_SENT_DIRECTLY, so finish() was never reached. The SKIPPED bullet below already branches on exactly this. Separately, 1807 and 1823-1824 scope the rule to "reported via message", but the code checks _terminal_state is not None, which also catches SERVER_TERMINATED set at 1778 by _send_heartbeat_if_needed with no subprocess message.

test_supervisor.py:4208 and :4245 say _exit_code = 1 comes "from a SIGTERM kill", but signal deaths are negative (1279/1285 compare against -signal), and the overtime path ends at -9 since the child's SIGTERM handler at task_runner.py:1563 returns and force=True escalates. A positive 1 with a terminal state already set is reachable via finalize() at 2460, inside the try whose except Exception exits 1 at 2475.

The body undersells the fix by one case. TaskState(FAILED) with should_retry=True is reachable via task_runner.py:1648, 1667 and 1885, all of which emit FAILED regardless of retry eligibility. On base a later non-zero exit turned that into UP_FOR_RETRY, so nothing was written and the row stayed running.

A subprocess can report a terminal state via message (SucceedTask,
TaskState(SKIPPED), etc.) and then exit with a genuinely non-zero code
afterward -- e.g. an OOM kill, or `_handle_process_overtime_if_needed()`
sending SIGTERM once `_terminal_state` is already set. `final_state`
previously derived its result from the exit code whenever one was
observed, ignoring an already-confirmed terminal state.

The clearest reachable case is SUCCESS: `SucceedTask` writes the row
directly and clears the pending-message slot on success, so a later
non-zero exit code re-triggers `update_task_state_if_needed()` ->
`.finish()` against a row `succeed()` already wrote, which 409s.

A terminal state can also be set directly by the supervisor rather
than reported by the subprocess (SERVER_TERMINATED via
`_send_heartbeat_if_needed`), so the precedence applies to any
already-set `_terminal_state`, not only ones delivered by message --
this subsumes needing SERVER_TERMINATED/UP_FOR_RETRY called out as
special cases, since every terminal state now gets the same
precedence.

Rebased over apache#73249/apache#73253's rework of pending-message delivery:
those changes mean TaskState-reported outcomes (SKIPPED, FAILED,
etc.) are now normally dispatched using the message's own state
directly, without consulting `final_state` at all while a message is
pending, so this precedence is defense-in-depth for that path rather
than the reachable bug it was in the pre-rework code -- the SUCCESS
case above remains the live, reachable one this fixes.

Confirmed via live instrumentation against a reproduced crash, see
apache#65708 (comment).
@seanmuth
seanmuth force-pushed the seanmuth/fix-ti-terminal-state-precedence branch from 86a3f75 to 617d920 Compare September 25, 2026 17:39
@seanmuth

seanmuth commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor Author

Addressed the three optional nits from the approval review, and rebased over #73249/#73253 in the process, worth a second look on the latter given your approval predates those merges.

The three nits:

  • The SUCCESS bullet now correctly scopes the 409 claim to should_retry=False; with retries enabled, the pre-fix code resolves to UP_FOR_RETRY (in STATES_SENT_DIRECTLY), so no write happens and no 409 follows, but the state is still wrong (UP_FOR_RETRY instead of SUCCESS). Also reworded "reported via message" to cover SERVER_TERMINATED, which the supervisor sets directly with no subprocess message at all.
  • The test comments ("genuinely observed, e.g. from a SIGTERM kill") were wrong for a positive exit code -- signal deaths are negative. Corrected to cite the actually-reachable positive-exit-code path: an exception raised inside finalize()'s callbacks (task_runner.py's outer except Exception ... sys.exit(1)).
  • The PR body now mentions the third case: TaskState(FAILED) with should_retry=True, reachable via task_runner.py regardless of retry eligibility, where the pre-fix code loses the FAILED report to a spurious UP_FOR_RETRY.

The rebase: #73249 restructured how TaskState-reported outcomes get dispatched -- they now go through update_task_state_if_needed()'s pending-message path using the message's own state directly, never consulting final_state while a message is pending. That makes the SKIPPED-case fix here defense-in-depth rather than the reachable bug it was when you reviewed it; the SUCCESS case remains fully live and reachable either way. #73253 also added its own narrower precedence special-case for SERVER_TERMINATED/UP_FOR_RETRY directly in final_state; the general _terminal_state is not None check here subsumes both under one rule, so I removed the narrower tuple rather than keeping both. Full diff against current main is now just the final_state property itself, nothing else touched.

294 passed (up from 211 pre-rebase, picking up #73249/#73253's own new tests), 1 skipped (platform-only), mypy clean.


Drafted-by: Claude Sonnet 5 (reviewed by @seanmuth)

This branch has not been deployed

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Airflow 3.2 Worker Failures on TaskInstance Finish (HTTP 409 invalid_state)

2 participants