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

shahar1 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 5523443d98c Set end date and duration on tasks skipped by a Dag run 
timeout (#74128)
5523443d98c is described below

commit 5523443d98c0aa9fabc74a40dd8f5d3c87317c90
Author: Lucas Pimentel <[email protected]>
AuthorDate: Sat Oct 3 04:49:01 2026 -0300

    Set end date and duration on tasks skipped by a Dag run timeout (#74128)
    
    * Test that a Dag-run timeout skips a running task with an end date
    
    When a Dag run hits dagrun_timeout, the scheduler skips the task that is
    still running but never sets its end_date or duration, so the elapsed time
    shown for it keeps growing. This test fails until the timeout path sets
    both.
    
    * Set end date and duration on tasks skipped by a Dag-run timeout
    
    The timeout branch of SchedulerJobRunner._schedule_dag_run assigned
    SKIPPED directly, so a task that was running kept its start_date with no
    end_date or duration, and the elapsed time shown for it kept growing.
    Use TaskInstance.set_state, which every other skip path goes through: it
    sets end_date and duration, and gives a task that never started the same
    start and end date.
---
 .../src/airflow/jobs/scheduler_job_runner.py       |  3 +-
 airflow-core/tests/unit/jobs/test_scheduler_job.py | 64 ++++++++++++++++++++++
 2 files changed, 65 insertions(+), 2 deletions(-)

diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py 
b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
index c727a93182c..5e07bf3cf70 100644
--- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
+++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
@@ -3152,8 +3152,7 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
                 default=None,
             )
             for task_instance in unfinished_task_instances:
-                task_instance.state = TaskInstanceState.SKIPPED
-                session.merge(task_instance)
+                task_instance.set_state(TaskInstanceState.SKIPPED, 
session=session)
             session.flush()
             self.log.info("Run %s of %s has timed-out", dag_run.run_id, 
dag_run.dag_id)
 
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py 
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index f0c92158e12..08d9588b8bc 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -4549,6 +4549,70 @@ class TestSchedulerJob:
         session.rollback()
         session.close()
 
+    def 
test_dagrun_timeout_skips_running_task_with_end_date_and_duration(self, 
dag_maker):
+        """A task still running when its Dag run times out must be skipped 
with an end_date and duration."""
+        session = settings.Session()
+        with dag_maker(
+            dag_id="test_scheduler_dagrun_timeout_running_task",
+            dagrun_timeout=datetime.timedelta(seconds=60),
+            session=session,
+        ):
+            EmptyOperator(task_id="dummy")
+
+        now = timezone.utcnow().replace(microsecond=0)
+        dr = dag_maker.create_dagrun(start_date=now - 
datetime.timedelta(days=1))
+        ti = dr.get_task_instance("dummy", session=session)
+        ti.state = TaskInstanceState.RUNNING
+        ti.start_date = now - datetime.timedelta(minutes=10)
+        session.flush()
+
+        scheduler_job = Job()
+        self.job_runner = SchedulerJobRunner(job=scheduler_job)
+
+        with time_machine.travel(now, tick=False):
+            self.job_runner._schedule_dag_run(dr, session)
+        session.flush()
+
+        session.refresh(ti)
+        assert ti.state == TaskInstanceState.SKIPPED
+        assert ti.end_date == now
+        assert ti.duration == 600.0
+
+        session.rollback()
+        session.close()
+
+    def test_dagrun_timeout_skips_unstarted_task_with_start_and_end_date(self, 
dag_maker):
+        """A task that never started gets start_date == end_date and a zero 
duration, like other skips."""
+        session = settings.Session()
+        with dag_maker(
+            dag_id="test_scheduler_dagrun_timeout_unstarted_task",
+            dagrun_timeout=datetime.timedelta(seconds=60),
+            session=session,
+        ):
+            EmptyOperator(task_id="dummy")
+
+        now = timezone.utcnow().replace(microsecond=0)
+        dr = dag_maker.create_dagrun(start_date=now - 
datetime.timedelta(days=1))
+        ti = dr.get_task_instance("dummy", session=session)
+        ti.state = None
+        session.flush()
+
+        scheduler_job = Job()
+        self.job_runner = SchedulerJobRunner(job=scheduler_job)
+
+        with time_machine.travel(now, tick=False):
+            self.job_runner._schedule_dag_run(dr, session)
+        session.flush()
+
+        session.refresh(ti)
+        assert ti.state == TaskInstanceState.SKIPPED
+        assert ti.start_date == now
+        assert ti.end_date == now
+        assert ti.duration == 0.0
+
+        session.rollback()
+        session.close()
+
     @mock.patch("airflow._shared.observability.metrics.stats._get_backend")
     def test_dagrun_timeout_duration_metric_has_run_type(self, 
mock_get_backend, dag_maker):
         """

Reply via email to