kaxil opened a new pull request, #73454:
URL: https://github.com/apache/airflow/pull/73454

   When a triggerer's heartbeat lapses for longer than 
`triggerer_health_check_threshold`, `Trigger.assign_unassigned` hands its 
triggers to another triggerer. The old triggerer is often still alive: the 
runner-silence watchdog deliberately skips heartbeats when the runner 
subprocess stops responding, and a long GC pause, a stopped container or a 
network partition do the same. When it recovers, 
`TriggerRunnerSupervisor.update_triggers` sees those triggers missing from its 
assignment list and queues them for cancellation with the user-action message, 
so `run_trigger` calls `on_kill()` on every one of them. Nothing in that path 
distinguishes "the row was deleted because a user cleared the task" from "the 
row was reassigned".
   
   For any trigger that cancels remote work from `on_kill()` this cancels a job 
the new owner has just started polling. The batch triggers in the Anthropic and 
Common AI providers, and the EMR, Redshift Data, Databricks, Dataproc and 
BigQuery triggers all do. The `BaseTrigger.on_kill` docstring and the deferring 
docs both promise it does not fire on redistribution, so this is a contract 
violation, present since `on_kill()` shipped in 3.3.0 (#65590). The 
`cancel_triggers` docstring claimed redistribution "goes through a separate 
path"; there was no such path.
   
   ## Solution
   
   Before queueing a cancel, the supervisor asks the DB which of the departed 
ids still have a trigger row and who owns them now 
(`Trigger.fetch_assignments`, behind an overridable `fetch_trigger_assignments` 
hook that owns its own session so `update_triggers` stays DB-free for 
subclasses without a metadata DB). A row that survives after leaving our list 
was reassigned, so those ids go to a new `releasing_triggers` set and reach the 
runner as `TriggerStateSync.to_release`; the runner cancels those asyncio tasks 
with a distinct sentinel, and `run_trigger` drops the trigger without 
`on_kill()`. Only a row that is gone, which is what a clear or 
mark-success/failed produces, takes the existing user-action path. The 
classification runs once per departed id, since ids already queued are excluded 
from the next loop's difference. The new field defaults to an empty set, so 
existing `TriggerStateSync` producers are unaffected.
   
   **Why in the supervisor rather than the runner.** The runner subprocess has 
no DB access; the supervisor already owns the assignment query and the per-loop 
`create_session()` in `build_trigger_workloads`, so the one extra `SELECT id 
... WHERE id IN (...)` sits next to it, and the runner keeps its existing 
message-based contract.
   
   ## Reproduction
   
   Two triggerers started with `--queues afl297` and 
`triggerer.queues_enabled=True` (so the standalone triggerer stays out of it), 
a task on that queue deferring to a trigger whose `on_kill()` appends a line to 
a file, then `SIGSTOP` on the owning triggerer and its runner child until the 
trigger row moved, then `SIGCONT`.
   
   | | Old owner after `SIGCONT` | `on_kill()` marker |
   |---|---|---|
   | main | `Trigger cancelled by user action, invoking on_kill` | written by 
the old owner's runner while the row already belonged to the new owner and the 
task was still `deferred` |
   | this PR | `Triggers were reassigned to another triggerer, releasing them 
without on_kill new_owners={3: 8}` and, in the runner, `Trigger reassigned to 
another triggerer, dropping it without on_kill` | not written |
   
   With this PR, marking the deferred task failed through the REST API still 
produced the marker from the new owner's runner within seconds, so the 
user-action path is intact.
   
   ## Known gap
   
   Row existence proves the row was reassigned, not that the new owner has 
started it. If a user clears the task in the window after the old owner 
classified the trigger as reassigned and before the new owner fetched its 
workload, neither side calls `on_kill()`: the old owner released, and the new 
owner skips a trigger whose task instance is gone. This is the same outcome 
main has when the old owner has actually died, so `on_kill()` on a user action 
was already best-effort while a failover is in flight. The fix errs toward 
leaving remote work running rather than cancelling work another triggerer is 
polling.
   
   ---
   
   * 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