This is an automated email from the ASF dual-hosted git repository.

hussein-awala pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 94a677e8f87 Keep callback triggers assigned to healthy triggerers 
(#72848)
94a677e8f87 is described below

commit 94a677e8f877eea59589aa77022f106a3edcf340
Author: Hussein Awala <[email protected]>
AuthorDate: Mon Oct 5 11:50:07 2026 +0200

    Keep callback triggers assigned to healthy triggerers (#72848)
    
    Callback triggers could be repeatedly claimed by another healthy triggerer, 
causing duplicate callback execution and consuming capacity intended for 
waiting triggers.
---
 airflow-core/src/airflow/models/trigger.py     |  1 +
 airflow-core/tests/unit/models/test_trigger.py | 45 ++++++++++++++++++++++++++
 2 files changed, 46 insertions(+)

diff --git a/airflow-core/src/airflow/models/trigger.py 
b/airflow-core/src/airflow/models/trigger.py
index 6a3e704fe79..aae270756be 100644
--- a/airflow-core/src/airflow/models/trigger.py
+++ b/airflow-core/src/airflow/models/trigger.py
@@ -462,6 +462,7 @@ class Trigger(Base):
             # Callback triggers
             select(cls.id)
             .join(Callback, isouter=False)
+            .where(or_(cls.triggerer_id.is_(None), 
cls.triggerer_id.not_in(alive_triggerer_ids)))
             .order_by(Callback.priority_weight.desc(), cls.created_date),
             # Task Instance triggers
             select(cls.id)
diff --git a/airflow-core/tests/unit/models/test_trigger.py 
b/airflow-core/tests/unit/models/test_trigger.py
index feb5c7aba6b..80961acb4d9 100644
--- a/airflow-core/tests/unit/models/test_trigger.py
+++ b/airflow-core/tests/unit/models/test_trigger.py
@@ -667,6 +667,51 @@ def test_assign_unassigned(session, create_triggerer, 
create_trigger, use_queues
         )
 
 
[email protected]("queue", [None, "callbacks"])
+def test_assign_unassigned_callbacks_preserves_healthy_owners(session, 
create_triggerer, time_machine, queue):
+    now = timezone.datetime(2026, 1, 1)
+    time_machine.move_to(now, tick=False)
+    queues = {queue} if queue else None
+    healthy_owner = create_triggerer(session, State.RUNNING, 
latest_heartbeat=now)
+    claiming_triggerer = create_triggerer(session, State.RUNNING, 
latest_heartbeat=now)
+    stale_owner = create_triggerer(
+        session, State.RUNNING, latest_heartbeat=now - 
datetime.timedelta(seconds=31)
+    )
+    finished_owner = create_triggerer(session, State.SUCCESS, 
latest_heartbeat=now, end_date=now)
+    session.flush()
+
+    expected_owners = {}
+    for owner, priority in (
+        (healthy_owner, 10),
+        (claiming_triggerer, 10),
+        (stale_owner, 1),
+        (finished_owner, 1),
+        (None, 1),
+    ):
+        callback = TriggererCallback(
+            callback_def=AsyncCallback("asyncio.sleep", kwargs={"delay": 60}, 
queue=queue),
+            priority_weight=priority,
+        )
+        callback.queue(session=session)
+        callback.trigger.triggerer_id = owner.id if owner else None
+        session.add(callback)
+        session.flush()
+        expected_owners[callback.trigger.id] = (
+            healthy_owner.id if owner is healthy_owner else 
claiming_triggerer.id
+        )
+    session.commit()
+
+    for triggerer in (claiming_triggerer, healthy_owner):
+        Trigger.assign_unassigned(
+            triggerer.id,
+            capacity=4,
+            health_check_threshold=30,
+            queues=queues,
+        )
+        session.expire_all()
+        assert dict(session.execute(select(Trigger.id, 
Trigger.triggerer_id)).all()) == expected_owners
+
+
 @pytest.mark.need_serialized_dag
 @conf_vars({("triggerer", "queues_enabled"): "True"})
 def test_assign_unassigned_with_qeueus(session, create_triggerer, 
create_trigger) -> None:

Reply via email to