kaxil commented on code in PR #73454:
URL: https://github.com/apache/airflow/pull/73454#discussion_r4086161925


##########
airflow-core/src/airflow/jobs/triggerer_job_runner.py:
##########
@@ -995,8 +1021,22 @@ def update_triggers(self, requested_trigger_ids: 
set[int]):
             self.creating_triggers.extend(workloads_to_create)
 
         if cancel_trigger_ids:
-            # Enqueue orphaned triggers for cancellation
-            self.cancelling_triggers.update(cancel_trigger_ids)
+            # A trigger leaves our assignment list for one of two reasons, and 
only the DB can tell
+            # them apart:
+            #   - The row is gone: the task instance is no longer deferred on 
it, because a user
+            #     cleared or marked it, or the scheduler timed the deferral 
out. The runner cancels
+            #     with the user-action message and run_trigger() decides 
whether on_kill() applies.
+            #   - The row still exists but points at another triggerer: our 
heartbeat lapsed and
+            #     Trigger.assign_unassigned handed the trigger over. The new 
owner is polling the
+            #     remote work, so the runner must drop it locally *without* 
on_kill().

Review Comment:
   Done in fe9fd85a2b.



##########
airflow-core/src/airflow/jobs/triggerer_job_runner.py:
##########
@@ -513,6 +519,10 @@ class TriggerRunnerSupervisor(WatchedSubprocess):
     # FinishedTriggers message
     cancelling_triggers: set[int] = attrs.field(factory=set, init=False)
 
+    # Like cancelling_triggers, but for triggers reassigned to another 
triggerer (see update_triggers);
+    # the async process drops these without invoking on_kill()

Review Comment:
   Removed in fe9fd85a2b.



##########
airflow-core/src/airflow/jobs/triggerer_job_runner.py:
##########
@@ -1442,19 +1486,21 @@ async def create_triggers(self):
 
     async def cancel_triggers(self):
         """
-        Drain the to_cancel queue and ensure all triggers that are not in the 
DB are cancelled.
+        Drain the to_cancel and to_release queues and cancel the matching 
trigger tasks.
 
-        This allows the cleanup job to delete them.
-        Passes "user-action" as the cancel message so that run_trigger() knows 
to invoke
-        on_kill(). Triggers in this queue are always removed because this is 
the path in which
-        the user performed some action on the task. Trigger redistribution 
goes through a separate
-        path.
+        The cancel message tells run_trigger() whether on_kill() applies: 
``to_cancel`` holds
+        triggers whose row is gone, ``to_release`` holds triggers reassigned 
to another
+        triggerer. See TriggerRunnerSupervisor.update_triggers for how the 
split is made.

Review Comment:
   Done in fe9fd85a2b. I kept it on one line but swapped "whether" for "if" so 
it fits within the 110-character limit.



##########
airflow-core/tests/unit/jobs/test_triggerer_job.py:
##########
@@ -2744,6 +2804,109 @@ def 
test_update_triggers_skips_when_ti_has_no_dag_version(session, supervisor_bu
     supervisor.stdin.write.assert_not_called()
 
 
+def _alive_triggerer_job(session) -> Job:
+    other_job = Job(job_type="TriggererJob")
+    other_job.latest_heartbeat = timezone.utcnow()
+    session.add(other_job)
+    session.flush()
+    return other_job
+
+
+def 
test_load_triggers_releases_reassigned_trigger_without_user_action_cancel(session,
 supervisor_builder):
+    """
+    A trigger whose row moved to another triggerer must be dropped locally, 
not queued
+    as a user-action cancel (which would fire ``on_kill()`` and cancel remote 
work the new
+    owner is still polling).
+    """
+    trigger = TimeDeltaTrigger(datetime.timedelta(days=7))
+    _, _, trigger_orm, _ = create_trigger_in_db(session, trigger)
+    supervisor = supervisor_builder()
+    supervisor.running_triggers = {trigger_orm.id}
+
+    # Our heartbeat lapsed and another (alive) triggerer took the trigger over.
+    other_job = _alive_triggerer_job(session)
+    session.execute(update(Trigger).where(Trigger.id == 
trigger_orm.id).values(triggerer_id=other_job.id))
+    session.flush()
+
+    supervisor.load_triggers()
+
+    assert supervisor.cancelling_triggers == set()
+    assert supervisor.releasing_triggers == {trigger_orm.id}
+
+
+def test_load_triggers_cancels_deleted_trigger_as_user_action(session, 
supervisor_builder):
+    """A trigger whose row is gone (task left the deferred state) takes the 
user-action path."""
+    trigger = TimeDeltaTrigger(datetime.timedelta(days=7))
+    _, _, trigger_orm, task_instance = create_trigger_in_db(session, trigger)
+    supervisor = supervisor_builder()
+    supervisor.running_triggers = {trigger_orm.id}
+
+    session.execute(
+        update(TaskInstance)
+        .where(TaskInstance.id == task_instance.id)
+        .values(trigger_id=None, state="scheduled")
+    )
+    session.execute(delete(Trigger).where(Trigger.id == trigger_orm.id))
+    session.flush()
+
+    supervisor.load_triggers()
+
+    assert supervisor.cancelling_triggers == {trigger_orm.id}
+    assert supervisor.releasing_triggers == set()
+
+
+def 
test_update_triggers_splits_cancel_set_by_row_existence(supervisor_builder, 
mocker):
+    supervisor = supervisor_builder()
+    supervisor.running_triggers = {1, 2, 3}
+    fetch_assignments = mocker.patch.object(
+        TriggerRunnerSupervisor, "fetch_trigger_assignments", autospec=True, 
return_value={2: 99}
+    )
+
+    supervisor.update_triggers({3})
+
+    fetch_assignments.assert_called_once_with(supervisor, {1, 2})
+    assert supervisor.releasing_triggers == {2}
+    assert supervisor.cancelling_triggers == {1}
+
+    # Already-classified triggers are not re-queried on the next loop.
+    supervisor.update_triggers({3})
+    fetch_assignments.assert_called_once()
+
+
+def 
test_update_triggers_does_not_query_when_nothing_to_cancel(supervisor_builder, 
mocker):
+    supervisor = supervisor_builder()
+    supervisor.running_triggers = {1}
+    fetch_assignments = mocker.patch.object(
+        TriggerRunnerSupervisor, "fetch_trigger_assignments", autospec=True
+    )
+
+    supervisor.update_triggers({1})
+
+    fetch_assignments.assert_not_called()

Review Comment:
   Agreed, removed in fe9fd85a2b.



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