1fanwang opened a new pull request, #73816:
URL: https://github.com/apache/airflow/pull/73816
<!-- SPDX-License-Identifier: Apache-2.0
https://www.apache.org/licenses/LICENSE-2.0 -->
When infrastructure stops a task worker, the worker cannot report an
exception. A Kubernetes preemption, an ECS Spot interruption or a lost Celery
worker all reach failure listeners, logs and metrics as a generic failure, so
integrations that route alerts or decide on cleanup have to parse
backend-specific text to tell an infrastructure interruption from an
application error.
This composite POC for [AIP-97 Task Failure
Classification](https://cwiki.apache.org/confluence/x/0ITMFw) carries an
optional failure category and a short diagnostic reason from the executor into
failure handling: the failure listener event, the scheduler failure log and a
bounded `failure_kind` label on the existing failure counters. Kubernetes
reports `infra` only for recognized disruption evidence such as
`DeletionByTaintManager`. The ECS executor preserves the task stop code and
reports `SpotInterruption` as infrastructure; other codes stay diagnostic
reasons without a category. The Celery executor preserves `WorkerLost` as a
reason with no category, because a lost worker alone does not establish the
cause. Anything ambiguous stays unclassified.
It grants no retries and adds no configuration, database fields or Dag
callback context, and clear behavior is unchanged. The separate [AIP-122
Infrastructure Retries](https://cwiki.apache.org/confluence/x/-pXwGg) policy
builds on this contract and has its own
[demonstration](https://github.com/1fanwang/airflow/pull/16).
## Testing Done
The branch stacks three sources that each ran end to end. Before merging
main, the composite head `6032b7ce28` was byte-identical to every tested source
for the files that source owns; this check printed nothing:
```bash
git diff --stat 6f7fa704122c6263be59eba6af4bb5d4264d1997 6032b7ce28 --
airflow-core shared task-sdk providers/cncf
git diff --stat fbd551cf76da7372317955a58caf880daef2a891 6032b7ce28 --
providers/celery
git diff --stat 5819598225a7392d1ba274b07666b9806db8f2ef 6032b7ce28 --
providers/amazon
```
The merge with main (`1d916545b2`) then resolved conflicts with the attempt
UUID and try number alignment from https://github.com/apache/airflow/pull/73554
in five files: `notify_failure()` also carries `failure_kind` and `reason`,
`_finalize_task_failure()` keeps the new `log` parameter next to
`failure_kind`, the scheduler event loop keeps both sides' imports, and the two
affected tests take the upstream fixtures. The end-to-end runs below predate
that merge; the new upstream listener test recorder now also accepts the two
cause keywords.
Every run used the real scheduler, Dag processor, authenticated API server
and Task SDK workers against PostgreSQL 16, with a UDP receiver capturing
emitted metrics. Each run compares a baseline without the change to the
candidate. Task rows arise through the scheduler and API; no fixture inserts
rows or mocks the executor client.
### Kubernetes and core classification
The [complete executed fixture and
receipts](https://gist.github.com/1fanwang/83c99b2db3759f0f776f6e096d39321c/984e9a4b57629e5599a7712d4f04b4ba2a61f689)
compare [baseline
b6538c32](https://github.com/apache/airflow/commit/b6538c32c8b1cdcb7b818e174a8b09d69ae18ad6)
with [candidate
6f7fa704](https://github.com/apache/airflow/commit/6f7fa704122c6263be59eba6af4bb5d4264d1997)
on a local kind cluster. After the fixture's source and image setup:
```bash
node orchestrate.mjs baseline 02
node orchestrate.mjs candidate 01
```
The candidate ran this inside Breeze; the fixture includes the full
controller and executed shell body:
```bash
export AIP97_PROJECT=aip97-c97-775ee60b-candidate-01
AIP97_DATABASE=aip97_c97_775ee60b_candidate_01
AIP97_IMAGE=aip97-classification-775ee60b-candidate:source
export AIP97_SOURCE_COMMIT=6f7fa704122c6263be59eba6af4bb5d4264d1997
AIP97_CLASSIFICATION=1
python /recipe/run_live.py
```
The controller waited for a real task-body marker and a later heartbeat
before applying a `NoExecute` taint. On the baseline the scheduler listener had
no cause:
```text
{"utc": "2026-09-27T03:05:22.367157+00:00", "origin": "scheduler", "dag_id":
"aip97_classification_infra", "task_id": "probe", "run_id":
"manual__infra_8704c008", "try_number": 1, "kind": null, "reason": null}
```
The candidate's raw observations show the accepted start and heartbeat, the
cause and metric, and a failed task with no additional retry:
```text
{"timestamp":"2026-09-27T03:13:03.049813Z","level":"info","event":"request
finished","method":"PATCH","path":"/execution/task-instances/01a0e0d9-ba26-7a5e-95f8-78e5d47017f7/run","query":"","status_code":200,"duration_us":86053,"client_addr":"172.20.0.2:54544","logger":"http.access","filename":"http_access_log.py","lineno":130}
{"timestamp":"2026-09-27T03:13:04.387600Z","level":"info","event":"request
finished","method":"PUT","path":"/execution/task-instances/01a0e0d9-ba26-7a5e-95f8-78e5d47017f7/heartbeat","query":"","status_code":204,"duration_us":40009,"client_addr":"172.20.0.2:54544","logger":"http.access","filename":"http_access_log.py","lineno":130}
{"utc": "2026-09-27T03:13:35.854655+00:00", "origin": "scheduler", "dag_id":
"aip97_classification_infra", "task_id": "probe", "run_id":
"manual__infra_95009922", "try_number": 1, "kind": "infra", "reason":
"DeletionByTaintManager"}
aip97.ti_failures:1|c|#dag_id:aip97_classification_infra,run_type:manual,task_id:probe,failure_kind:infra|c:in-210068
{"utc": "2026-09-27T03:13:35.948122+00:00", "ti": {"id":
"01a0e0d9-ba26-7a5e-95f8-78e5d47017f7", "dag_id": "aip97_classification_infra",
"task_id": "probe", "run_id": "manual__infra_95009922", "state": "failed",
"try_number": 1, "max_tries": 0, "queued_by_job_id": 2, "start_date":
"2026-09-27 03:12:42.719757+00:00", "end_date": "2026-09-27
03:13:35.853122+00:00", "last_heartbeat_at": "2026-09-27
03:13:34.096197+00:00", "pid": 14}}
```
A workload-storage control wrote 8 MiB into an `emptyDir` limited to 1 MiB.
Its eviction stayed unclassified:
```text
{"utc": "2026-09-27T03:14:06.553046+00:00", "origin": "scheduler", "dag_id":
"aip97_classification_storage", "task_id": "probe", "run_id":
"manual__storage_4c49eefd", "try_number": 1, "kind": null, "reason": "Evicted"}
```
Application, timeout and manual-failure controls exercised their reporting
paths, and an ordinary retry plus a public-API clear kept callback context
unchanged:
```text
{"callback": "retry", "dag_id": "aip97_classification_retry", "run_id":
"manual__retry_711428d1", "try_number": 1, "state": "up_for_retry",
"cause_in_context": false}
{"callback": "success", "dag_id": "aip97_classification_retry", "run_id":
"manual__retry_711428d1", "try_number": 2, "state": "success",
"cause_in_context": false}
```
The receipt covers 12 scenarios and 16 authenticated task attempts with
matching retry, clear, callback and schema observations. The taint-manager case
exited 137 after the default 30-second grace period, so it does not establish
cooperative SIGTERM behavior, and this run does not cover scheduler preemption,
standalone OOM or SIGKILL, or missed heartbeats.
### Celery worker loss
The [complete executed fixture, raw observations and independent
verification](https://gist.github.com/1fanwang/1fd60e00f3d18b0e7be92aa478196cf6/f9d826d8dccbe3ab99d49da8a704c82a77de3b7e)
compare baseline `6f7fa704122c6263be59eba6af4bb5d4264d1997` with candidate
`fbd551cf76da7372317955a58caf880daef2a891`. Real Dags ran through the API, Dag
processor, scheduler, Celery prefork worker and Task SDK against PostgreSQL, on
both Redis/Redis and RabbitMQ/PostgreSQL, across worker loss with zero or one
ordinary retry, an application error containing `signal 9 (SIGKILL)`, and
success. Every kill followed an accepted authenticated start and a later
heartbeat.
```bash
env -u DOCKER_HOST DOCKER_CONTEXT=colima DOCKER_DEFAULT_PLATFORM=linux/amd64
\
breeze run --python 3.12 --platform linux/amd64 --backend none \
--project-name aip97-celery-runtime --skip-image-upgrade-check --answer no
\
env GIT_DIR="$(git rev-parse --absolute-git-dir)"
GIT_WORK_TREE=/opt/airflow \
python dev/aip97_celery_runtime/run.py \
--candidate fbd551cf76da7372317955a58caf880daef2a891 \
--redis-image redis:7-alpine --postgres-image postgres:14 \
--rabbitmq-image
rabbitmq@sha256:47c388129d4fc820f399d60d4e358c03a91617a657be024be7321be3387de819
```
The candidate's Redis task started and heartbeated before its Celery child
was killed:
```json
{"timestamp":"2026-09-27T06:01:39.483397Z","level":"info","event":"request
finished","method":"PATCH","path":"/execution/task-instances/01a0e174-148e-7858-bb22-6fed40143d0e/run","query":"","status_code":200,"duration_us":95392,"client_addr":"127.0.0.1:44180","logger":"http.access","filename":"http_access_log.py","lineno":130}
{"timestamp":"2026-09-27T06:01:50.063215Z","level":"info","event":"request
finished","method":"PUT","path":"/execution/task-instances/01a0e174-148e-7858-bb22-6fed40143d0e/heartbeat","query":"","status_code":204,"duration_us":21405,"client_addr":"127.0.0.1:44180","logger":"http.access","filename":"http_access_log.py","lineno":130}
{"utc": "2026-09-27T06:01:52.965883+00:00", "event": "BACKEND_WORKER_LOST",
"killed_at": "2026-09-27T06:01:50.761581+00:00", "celery_task_id":
"6d25e9f7-f7f0-44a7-85df-b3cd569328f9", "worker_pid": 2071, "pid_created_at":
1790488860.9, "signal": 9, "backend_state": "FAILURE", "exception_type":
"WorkerLostError", "error": "Worker exited prematurely: signal 9 (SIGKILL) Job:
0."}
```
The same failure on RabbitMQ reached the scheduler listener without a reason
on the baseline and with the reason on the candidate; neither supplied a
failure kind:
```json
{"utc": "2026-09-27T05:43:46.273053+00:00", "origin": "scheduler", "dag_id":
"aip97_celery_lost_no_retry", "task_id": "probe", "run_id":
"manual__lost_no_retry_777984d0", "try_number": 1, "kind": null, "reason": null}
{"utc": "2026-09-27T05:56:00.772444+00:00", "origin": "scheduler", "dag_id":
"aip97_celery_lost_no_retry", "task_id": "probe", "run_id":
"manual__lost_no_retry_d0e3c460", "try_number": 1, "kind": null, "reason":
"WorkerLost"}
```
Raw candidate counters distinguish worker loss from the application control,
and the one-retry case recovered on attempt two with `max_tries=1` unchanged:
```text
aip97.ti_failures:1|c|#dag_id:aip97_celery_lost_no_retry,run_type:manual,task_id:probe,failure_kind:unclassified|c:in-249177
aip97.ti_failures:1|c|#dag_id:aip97_celery_application,task_id:probe,run_type:manual,failure_kind:application|c:in-249177
```
```json
{"utc": "2026-09-27T06:02:36.470670+00:00", "event": "CASE_FINISHED",
"phase": "candidate", "broker": "redis", "case": "lost_retry", "final": {"id":
"01a0e174-f4f5-76e5-8c3a-a927bee3635b", "dag_id": "aip97_celery_lost_retry",
"task_id": "probe", "run_id": "manual__lost_retry_2b920064", "state":
"success", "try_number": 2, "max_tries": 1, "queued_by_job_id": 2,
"start_date": "2026-09-27 06:02:19.502874+00:00", "end_date": "2026-09-27
06:02:30.246202+00:00", "last_heartbeat_at": "2026-09-27
06:02:30.053232+00:00", "pid": 2152, "hostname": "6f33a29f4c20",
"external_executor_id": "36deb26c-8bc9-4adf-8a0a-cbc0b5540243"}}
```
The independent verifier checked 20 task handshakes, all 31,478 captured
datagrams and 24 failure-counter increments, and confirmed all eight backing
containers were removed. Only signed tokens and broker credentials are redacted
in the published logs.
### ECS Spot interruption
This is real Airflow and boto HTTP against a local ECS substitute, not live
AWS. The [complete executed fixture, raw excerpts and source
hashes](https://gist.github.com/1fanwang/56a9a2e49c39d1eaf9b3cc41faaf0f32/27513757a44e7da543f9f895282c08682897107e)
compare the unchanged provider at `6f7fa704122c6263be59eba6af4bb5d4264d1997`
with the adapter at `5819598225a7392d1ba274b07666b9806db8f2ef`, using
PostgreSQL 16, Python 3.12 in Breeze and boto3/botocore 1.43.75. From a
checkout at the candidate commit:
```bash
FIXTURE=$(mktemp -d)
git clone --quiet
https://gist.github.com/56a9a2e49c39d1eaf9b3cc41faaf0f32.git "$FIXTURE"
git -C "$FIXTURE" checkout --quiet 27513757a44e7da543f9f895282c08682897107e
mkdir "$FIXTURE/recipe"
mv "$FIXTURE"/*.py "$FIXTURE/recipe/"
node "$FIXTURE/run.mjs" baseline 04
node "$FIXTURE/run.mjs" candidate 01
```
The same Spot-shaped response reaches the scheduler listener as:
```json
{"utc": "2026-09-27T05:51:12.336101+00:00", "origin": "scheduler", "dag_id":
"ecs_local_spot", "run_id": "manual__spot_c0ec9593", "try_number": 1, "kind":
null, "reason": null}
{"utc": "2026-09-27T06:03:32.673548+00:00", "origin": "scheduler", "dag_id":
"ecs_local_spot", "run_id": "manual__spot_90a11a22", "try_number": 1, "kind":
"infra", "reason": "SpotInterruption"}
```
Both phases keep identical terminal states and retry allowances. Missing and
other stop codes do not become infrastructure, including when the stopped
reason contains Spot-looking text, an ordinary retry succeeds on attempt two,
and a real deadline callback is interrupted after its authenticated start
without failing its still-running task. The independent verifier passed the
baseline and candidate phases with 17 authenticated task starts and later
heartbeats, raw failure counters, callback isolation, exact loaded source
hashes and cleanup. No AWS credentials or cloud resources were used, and no
failure-reason text was added to metric labels.
---
##### Was generative AI tooling used to co-author this PR?
- [X] Yes (please specify the tool below)
Generated-by: GitHub Copilot CLI (GPT-6) and Claude Code (Fable 5.1)
following [the
guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions)
---
* Read the **[Pull Request
Guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#pull-request-guidelines)**
for more information. Note: commit author/co-author name and email in commits
become permanently public when merged.
* For fundamental code changes, an Airflow Improvement Proposal
([AIP](https://cwiki.apache.org/confluence/display/AIRFLOW/Airflow+Improvement+Proposals))
is needed.
* When adding dependency, check compliance with the [ASF 3rd Party License
Policy](https://www.apache.org/legal/resolved.html#category-x).
* For significant user-facing changes create newsfragment:
`{pr_number}.significant.rst`, in
[airflow-core/newsfragments](https://github.com/apache/airflow/tree/main/airflow-core/newsfragments).
You can add this file in a follow-up commit after the PR is created so you
know the PR number.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]