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]