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


##########
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:
   ```suggestion
               # Only the DB tells the two cases apart: a gone row means the 
task left the
               # deferred state (user action), a surviving row means 
assign_unassigned handed
               # the trigger to a triggerer that is now polling the remote work.
   ```



##########
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:
   ```suggestion
   ```
   
   We can remove this ig, the attribute name and the block above say it



##########
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:
   `cancel_trigger_ids` is an empty set here. The query sits under `if 
cancel_trigger_ids:`. An empty set is falsy. So the test is asserting that 
Python skips a false if. That part cannot break.



##########
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:
   ```suggestion
   """Drain to_cancel and to_release, cancelling each task with the message 
that tells run_trigger() whether on_kill() applies."""
   ```
   
   shorter?



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