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]

Reply via email to