rino0601 opened a new issue, #71540:
URL: https://github.com/apache/airflow/issues/71540

   ### Under which category would you file this issue?
   
   Airflow Core
   
   ### Apache Airflow version
   
   3.3.0
   
   ### What happened and how to reproduce it?
   
   Deleting a single task instance row from a non-terminal dag run leaves that 
run permanently
   non-terminal. The row is never re-created, the run can never reach a 
terminal state, and it holds a
   `max_active_runs` slot indefinitely. On a DAG with `max_active_runs=1` this 
stops scheduling for
   that DAG completely.
   
   Crucially, the scheduler's own safety net — the "all tasks deadlocked" check 
that would otherwise
   fail the run and release the slot — is silently disabled for any run whose 
unfinished tasks carry a
   concurrency limit.
   
   Evidence below is marked **[observed]** (measured against a running 3.3.0 
instance) or
   **[code]** (read from the 3.3.0 source). We did not have server-log access, 
so no tracebacks.
   
   ### Reproduction
   
   ```python
   from datetime import datetime
   
   from airflow.sdk import dag, task
   
   
   @dag(
       dag_id="deleted_ti_strands_run",
       schedule="*/5 * * * *",
       start_date=datetime(2026, 1, 1),
       catchup=False,
       max_active_runs=1,
   )
   def deleted_ti_strands_run():
       @task
       def first():
           pass
   
       @task
       def second():
           pass
   
       # Any concurrency limit on a task that will remain unfinished is enough 
to
       # disable the deadlock check for this run (see "Why nothing recovers 
it", 3).
       @task(max_active_tis_per_dag=1)
       def third():
           pass
   
       first() >> second() >> third()
   
   
   deleted_ti_strands_run()
   ```
   
   1. Unpause the DAG and wait for a run in which `first` has succeeded and 
`second` has not started.
   2. Delete `second`'s task instance:
      ```
      DELETE 
/api/v2/dags/deleted_ti_strands_run/dagRuns/{run_id}/taskInstances/second?map_index=-1
      ```
   3. Observe:
      - `GET .../taskInstances/second` → `404`; the run now has one fewer task 
instance than the DAG
        defines.
      - `third` stays stateless (`null`) forever; the dag run stays `running` 
forever.
      - No `Task deadlock (no runnable tasks); marking run … failed` appears in 
the scheduler log.
      - No further dag run is ever created, because the stranded run holds the 
only
        `max_active_runs` slot.
   
   ### What we saw in production **[observed]**
   
   Four dag runs across three DAGs, each missing exactly one task instance row:
   
   | dag run | started | stranded for |
   |---|---|---|
   | `dag_a` (hourly) `04:00` | 04:00:00.673980Z | **4h31m** |
   | `dag_a` (hourly) `05:00` | 05:00:03.208903Z | 3h31m |
   | `dag_b` (`*/20`, `max_active_runs=1`) `05:00` | 05:14:48.517958Z | 3h16m |
   | `dag_c` (`*/10`, `max_active_runs=6`) `07:20` | 07:20:00.313212Z | 1h11m |
   
   None recovered on its own. All four recovered **simultaneously, within ~50 
seconds**, when the
   three DAG files were edited so their serialized representation changed — see 
"The only recovery
   path" below.
   
   `dag_b` is the clearest impact case: `*/20` schedule, `max_active_runs=1`, 
`catchup=False`. Between
   05:00Z and 09:00Z **exactly one dag run existed** (the stranded one). About 
**10 scheduled runs
   were never created**, and because `catchup=False` they are not backfilled — 
they are simply lost.
   
   Note on scope: this does not freeze *every* task. Branches that do not 
depend on the deleted task
   can still finish — in `dag_c` one teardown progressed to `success` on its 
own. The accurate
   statement is that **the chain through the deleted task stalls permanently 
and the run never reaches
   a terminal state**, so the slot is never released.
   
   `dagrun_timeout` is not a rescue: it is optional, and it was set to 30 hours 
on these DAGs
   **[observed]**.
   
   ### Why nothing recovers it
   
   There are three code paths that could restore the row or release the slot. 
All three are gated off.
   
   **1. Scheduler `verify_integrity` — gated on a DAG version change *and* on 
bundle versioning.**
   `airflow/jobs/scheduler_job_runner.py:2886`:
   
   ```python
   if not dag_run.bundle_version and not self._verify_integrity_if_dag_changed(
       dag_run=dag_run, session=session
   ):
   ```
   
   `_verify_integrity_if_dag_changed` (`:2923-2937`) returns early with
   `"DAG %s not changed structure, skipping dagrun.verify_integrity"` unless a 
new `DagVersion` exists.
   So an unchanged DAG never has its missing task instances restored. **[code]**
   
   Note the leading `not dag_run.bundle_version`: when a deployment uses bundle 
versioning, this call
   is skipped entirely, so even bumping the DAG version would not restore the 
row. Our runs had
   `bundle_version = None`, which is why the workaround below worked for us; we 
have **not** tested a
   bundle-versioned deployment, but by inspection it appears to have no 
recovery path at all. **[code]**
   
   **2. `clear_task_instances` `verify_integrity` — unreachable while the run 
is not terminal.**
   `airflow/models/taskinstance.py:438-449` calls `dr.verify_integrity(...)` 
only inside
   `if dr.state in State.finished_dr_states:` and only when 
`run_on_latest_version` is set. A
   `running` dag run takes the `elif run_on_latest_version:` branch, which just 
refreshes
   `created_dag_version_id`. So clearing cannot restore the row while the run 
is stranded — and the
   run cannot become terminal because it is stranded. **[code]**
   
   **3. Deadlock detection — disabled by any concurrency limit.** 
`airflow/models/dagrun.py:1280`:
   
   ```python
   elif unfinished.should_schedule and not are_runnable_tasks:
       self.log.error("Task deadlock (no runnable tasks); marking run %s 
failed", self)
       self.set_state(DagRunState.FAILED)
   ```
   
   This is exactly the situation — tasks are unfinished and none are runnable — 
so the run should be
   failed and its slot released. But `should_schedule` 
(`airflow/models/dagrun.py:1171-1185`) is:
   
   ```python
   return (
       bool(self.tis)
       and all(not getattr(t.task, "depends_on_past", False) for t in self.tis 
if t.task)
       and all(getattr(t.task, "max_active_tis_per_dag", None) is None for t in 
self.tis if t.task)
       and all(getattr(t.task, "max_active_tis_per_dagrun", None) is None for t 
in self.tis if t.task)
       and all(t.state not in (TaskInstanceState.DEFERRED, 
TaskInstanceState.AWAITING_INPUT) for t in self.tis)
   )
   ```
   
   Because every clause is `all(...)`, a **single** unfinished task carrying 
one of these attributes
   switches `should_schedule` off for the whole run. In one of our DAGs exactly 
one task out of nine
   carried `max_active_tis_per_dag=1`, and that alone was enough; the other two 
set it for every task
   via `default_args`. `max_active_tis_per_dag` is the 3.x successor to 2.x 
`task_concurrency`, so
   this is an ordinary operational setting, not an exotic one. **[code]**
   
   The ordering matters too. The deadlock branch also requires `not 
are_runnable_tasks`, so while
   independent branches can still make progress — in `dag_c` a teardown 
completed on its own — the
   condition is not yet met. The evaluation point only arrives once those 
branches have drained, and
   at that moment `should_schedule` is `False`, so the run is passed over 
silently.
   
   The outcome therefore bifurcates on DAG configuration:
   
   - **Some unfinished task carries a concurrency limit / `depends_on_past` / 
is deferred** → the run
     stays `running` indefinitely and holds its `max_active_runs` slot. Our 
four runs.
   - **Otherwise** → deadlock detection fires once independent branches drain, 
the run is marked
     `failed`, and the slot is released. A data gap remains where the deleted 
task should have run,
     but scheduling is not blocked.
   
   So a concurrency limit on any one task converts a self-limiting failure into 
an indefinite one.
   
   ### The only recovery path we found **[observed]**
   
   Editing the three DAG files so their serialized form changed (a `doc_md` 
change — Python comments
   are not serialized and have no effect) produced new `DagVersion` rows at 
08:30:54–08:31:04Z, and
   all four stranded runs were repaired by 08:31:10Z, with the 
previously-deleted tasks running. No
   database surgery, no dag run deletion, no per-run action — one deployment 
fixed all four.
   
   That this works is itself the proof of path 1's gate: recovery happened 
only, and immediately, when
   the DAG version changed.
   
   ### Attempts that did not work **[observed]**
   
   - `POST /api/v2/dags/{dag_id}/dagRuns/{run_id}/clear` with
     `{"dry_run": false, "only_failed": false, "run_on_latest_version": true}` 
on the stranded
     (`running`) run → **HTTP 500**. Nothing was committed (`clear_number` 
remained `None`). The same
     request with `"dry_run": true` returned `200` and correctly listed the 9 
surviving task
     instances. A 500 on a documented operation looks like a separate defect; 
we could not capture a
     traceback.
   - `PATCH /api/v2/dags/{dag_id}/dagRuns/{run_id}` with `{"state": "failed"}`, 
to force the run
     terminal so path 2 would become reachable → returned **`200` but the dag 
run state did not
     change**, and the response echoed `running`. Side effects were still 
committed: three pending
     non-teardown task instances were set to `skipped`.
     By inspection **[code]**, `set_dag_run_state_to_failed` in 
`airflow/api/common/mark_tasks.py`
     skips pending non-teardown TIs first and only then conditionally sets the 
run state:
     ```python
     if commit:
         for ti in pending_normal_tis:
             ti.set_state(TaskInstanceState.SKIPPED)
   
         # Mark the dag run to failed if there is no pending teardown (else 
this would not be scheduled later).
         if not any(dag.task_dict[ti.task_id].is_teardown for ti in 
(running_tis + pending_tis)):
             _set_dag_run_state(dag.dag_id, run_id, DagRunState.FAILED, session)
     ```
     With pending teardowns the guard holds, so the caller gets a `200`, a 
stale state echo, and a
     partially applied change. `state: "success"` has the same guard.
   
   Happy to split either of these into their own issues if preferred.
   
   ### What you think should happen instead?
   
   Any one of these would have prevented the outage. (2) is what an operator 
actually expects to
   happen, and (3) seems most clearly a bug.
   
   1. **The delete confirmation should describe the consequence.** The UI does 
warn before deleting —
      `"This will remove all metadata related to the {{type}}."` **[observed]** 
— but that describes
      *what is removed*, not *what it does to the run*. Nothing in it suggests 
that the run will never
      reach a terminal state, that it will hold a `max_active_runs` slot 
indefinitely, or that the DAG
      may stop scheduling altogether. An operator can read and accept that 
dialog and still have no
      way to anticipate an outage.
   
   2. **Re-create missing task instances for non-terminal runs.** This is what 
we expected: with the
      row gone, the scheduler would repopulate it on a later pass, the same way 
task instances are
      created for a new run. Instead the run is left in a state the scheduler 
will not resolve. The
      machinery already exists — `verify_integrity` does exactly this — it is 
just behind a gate that
      this situation never satisfies. Repairing a run whose task instance set 
does not match its own
      DAG version would not require the version to *change*.
   
   3. **Do not let concurrency limits disable deadlock detection.** 
`should_schedule` conflates "these
      tasks may be throttled" with "this run may still make progress". A run 
with unfinished tasks and
      no runnable tasks is deadlocked regardless of whether 
`max_active_tis_per_dag` is set; it should
      be failed so the `max_active_runs` slot is released.
   
   At minimum, a stranded run should not be able to consume a `max_active_runs` 
slot forever with no
   log line indicating why the DAG stopped scheduling.
   
   ### Operating System
   
   Linux (container)
   
   ### Deployment
   
   Other 3rd-party Helm chart
   
   ### Apache Airflow Provider(s)
   
   _No response_
   
   ### Versions of Apache Airflow Providers
   
   _No response_
   
   ### Official Helm Chart version
   
   Not Applicable
   
   ### Kubernetes Version
   
   Not Applicable
   
   ### Helm Chart configuration
   
   Not Applicable
   
   ### Docker Image customizations
   
   Not Applicable
   
   ### Anything else?
   
   **How to detect affected runs** — deadlocked runs are permanently 
non-terminal, so no time window
   is needed:
   
   ```
   GET /api/v2/dags/~/dagRuns?state=running&state=queued
   ```
   then, per run, compare the task ids from `/dags/{dag_id}/tasks` against 
those from
   `/dags/{dag_id}/dagRuns/{run_id}/taskInstances`; a non-empty difference is a 
missing row. Two
   caveats: a task instance with state `null` **exists** and is normal (not yet 
scheduled) — only a
   completely absent row counts; and comparing task ids will not detect a 
single deleted `map_index`
   of a mapped task. The query also finds only the stranded branch described 
above — a run that
   deadlock detection closed as `failed`, or one that closed as `success` 
because the deleted task was
   a leaf, is terminal and will not appear. Those leave a data gap without 
blocking scheduling.
   
   **Frequency:** four occurrences in a single day on one instance, all from an 
operator deleting a
   task instance where they meant to clear it. That figure counts the runs that 
were *blocking
   scheduling*, so it is a lower bound on how many task instances were actually 
deleted.
   
   ### Related
   
   - #48538 — a different failure when clearing a dag run (dynamically mapped 
tasks), also surfacing as
     a 5xx/409 from the clear endpoint.
   - #70056 — different cause, similar end state: dag runs that the scheduler 
silently stops
     progressing.
   
   ### Are you willing to submit PR?
   
   - [x] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
   


-- 
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