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):
"""