This is an automated email from the ASF dual-hosted git repository.
vatsrahul1001 pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/v3-3-test by this push:
new 70d8fa628c9 Fix Calendar view computing cron planned runs in UTC
instead of the Dag's timezone (#71243) (#72839)
70d8fa628c9 is described below
commit 70d8fa628c9647958b9738f0f18f762f0cac9ce9
Author: Rahul Vats <[email protected]>
AuthorDate: Thu Sep 10 08:10:30 2026 +0530
Fix Calendar view computing cron planned runs in UTC instead of the Dag's
timezone (#71243) (#72839)
* Fix Calendar planned runs using UTC instead of Dag timezone
CalendarService._calculate_cron_planned_runs() handed croniter a
UTC-tagged start_time, so cron fields were matched against UTC
wall-clock instead of the timetable's configured timezone. Step
planned instants with CronMixin._get_next() instead of iterating a
raw croniter, so the calendar matches the scheduler exactly,
including its DST gap/fold handling for every-hour crons.
Closes #71234
* Add newsfragment for #71243
* Remove newsfragment; not needed for small bugfixes
(cherry picked from commit 36a9a63ac08160772dd513605942eca45e5570ad)
# Conflicts:
#
airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_calendar.py
Co-authored-by: Sanghoon Kim / κΉμν
<[email protected]>
---
.../api_fastapi/core_api/services/ui/calendar.py | 26 ++---
.../core_api/routes/ui/test_calendar.py | 121 ++++++++++++++++++++-
2 files changed, 132 insertions(+), 15 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/services/ui/calendar.py
b/airflow-core/src/airflow/api_fastapi/core_api/services/ui/calendar.py
index 3c22472f734..373aca31088 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/services/ui/calendar.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/services/ui/calendar.py
@@ -17,14 +17,12 @@
from __future__ import annotations
import collections
-import itertools
-from collections.abc import Iterator, Sequence
+from collections.abc import Sequence
from datetime import datetime
from typing import Literal, cast
import sqlalchemy as sa
import structlog
-from croniter.croniter import croniter
from sqlalchemy.engine import Row
from sqlalchemy.orm import InstrumentedAttribute, Session
@@ -213,17 +211,17 @@ class CalendarService:
"""Calculate planned runs for cron-based timetables."""
dates: dict[datetime, int] = collections.Counter()
- dates_iter: Iterator[datetime | None] = croniter(
- cast("CronMixin", dag.timetable)._expression,
- start_time=last_data_interval.end,
- ret_type=datetime,
- )
-
- # Cap the iteration like _calculate_timetable_planned_runs does; a
high-frequency
- # expression (e.g. "* * * * *", or a seconds-resolution cron) would
otherwise take
- # hundreds of thousands of steps before hitting the year boundary.
- for dt in itertools.islice(dates_iter, self.MAX_PLANNED_RUNS):
- if dt is None or dt.year != year:
+ cron_timetable = cast("CronMixin", dag.timetable)
+ dt = last_data_interval.end
+
+ # Step with CronMixin._get_next so planned instants match the
scheduler exactly,
+ # including its DST gap/fold handling. Cap the iteration like
+ # _calculate_timetable_planned_runs does; a high-frequency expression
(e.g.
+ # "* * * * *", or a seconds-resolution cron) would otherwise take
hundreds of
+ # thousands of steps before hitting the year boundary.
+ for _ in range(self.MAX_PLANNED_RUNS):
+ dt = cron_timetable._get_next(dt)
+ if dt.year != year:
break
if dag.end_date and dt > dag.end_date:
break
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_calendar.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_calendar.py
index 5c218e599c2..da07ae89108 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_calendar.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/ui/test_calendar.py
@@ -25,7 +25,7 @@ from sqlalchemy.orm import Session
from airflow.api_fastapi.core_api.services.ui.calendar import CalendarService
from airflow.providers.standard.operators.empty import EmptyOperator
-from airflow.sdk import CronPartitionTimetable
+from airflow.sdk import CronPartitionTimetable, CronTriggerTimetable
from airflow.utils.session import NEW_SESSION, provide_session
from airflow.utils.state import DagRunState
@@ -187,6 +187,125 @@ class TestCalendar:
assert body == result
+class TestCalendarCronNonUTCTimezone:
+ """Planned runs for a cron timetable must be computed in the timetable's
own timezone, not UTC."""
+
+ DAG_NAME = "test_dag_non_utc_tz"
+
+ @pytest.fixture(autouse=True)
+ @provide_session
+ def setup_dag_runs(self, dag_maker, *, session: Session = NEW_SESSION) ->
None:
+ clear_db_runs()
+ clear_db_dags()
+ with dag_maker(
+ self.DAG_NAME,
+ schedule=CronTriggerTimetable("0 8 * * *", timezone="Asia/Seoul"),
+ start_date=datetime(2025, 1, 1),
+ catchup=True,
+ serialized=True,
+ session=session,
+ ):
+ EmptyOperator(task_id="test_task1")
+ dag_maker.create_dagrun(
+ run_id="run_1",
+ state=DagRunState.SUCCESS,
+ logical_date=pendulum.datetime(2025, 1, 1, 23, 0, 0, tz="UTC"),
+ )
+ dag_maker.sync_dagbag_to_db()
+
+ session.commit()
+
+ def teardown_method(self) -> None:
+ clear_db_runs()
+ clear_db_dags()
+
+ def test_planned_runs_use_timetable_timezone_not_utc(self, test_client):
+ response = test_client.get(f"/calendar/{self.DAG_NAME}",
params={"granularity": "hourly"})
+ assert response.status_code == 200
+ body = response.json()
+
+ planned = [r for r in body["dag_runs"] if r["state"] == "planned"]
+ # Daily 08:00 Asia/Seoul is 23:00Z the previous day; the last run's
data interval
+ # ends 2025-01-01T23:00Z, so planned runs are one per remaining day of
2025.
+ assert len(planned) == 364
+ assert min(r["date"] for r in planned) == "2025-01-02T23:00:00Z"
+ assert all(r["date"].endswith("T23:00:00Z") for r in planned), planned
+ assert all(r["count"] == 1 for r in planned)
+
+
+class CalendarEveryHourCronDstBase:
+ """Every-hour crons must plan exactly one run per UTC hour across a DST
transition, like the scheduler."""
+
+ DAG_NAME: str
+ START_DATE: datetime
+ LAST_RUN_UTC: pendulum.DateTime
+ EXPECTED_HOURS: list[str]
+
+ @pytest.fixture(autouse=True)
+ @provide_session
+ def setup_dag_runs(self, dag_maker, *, session: Session = NEW_SESSION) ->
None:
+ clear_db_runs()
+ clear_db_dags()
+ with dag_maker(
+ self.DAG_NAME,
+ schedule=CronTriggerTimetable("0 * * * *",
timezone="Europe/Zurich"),
+ start_date=self.START_DATE,
+ catchup=True,
+ serialized=True,
+ session=session,
+ ):
+ EmptyOperator(task_id="test_task1")
+ dag_maker.create_dagrun(
+ run_id="run_1",
+ state=DagRunState.SUCCESS,
+ logical_date=self.LAST_RUN_UTC,
+ )
+ dag_maker.sync_dagbag_to_db()
+
+ session.commit()
+
+ def teardown_method(self) -> None:
+ clear_db_runs()
+ clear_db_dags()
+
+ def test_one_planned_run_per_utc_hour(self, test_client):
+ response = test_client.get(f"/calendar/{self.DAG_NAME}",
params={"granularity": "hourly"})
+ assert response.status_code == 200
+
+ planned = {r["date"]: r["count"] for r in response.json()["dag_runs"]
if r["state"] == "planned"}
+ assert {h: planned.get(h) for h in self.EXPECTED_HOURS} ==
dict.fromkeys(self.EXPECTED_HOURS, 1)
+
+
+class TestCalendarEveryHourCronDstFold(CalendarEveryHourCronDstBase):
+ """Fall-back (2025-10-26 03:00 CEST -> 02:00 CET): the repeated hour keeps
its planned run."""
+
+ DAG_NAME = "test_dag_every_hour_dst_fold"
+ START_DATE = datetime(2025, 10, 1)
+ LAST_RUN_UTC = pendulum.datetime(2025, 10, 25, 22, 0, 0, tz="UTC")
+ EXPECTED_HOURS = [
+ "2025-10-25T23:00:00Z",
+ "2025-10-26T00:00:00Z",
+ "2025-10-26T01:00:00Z",
+ "2025-10-26T02:00:00Z",
+ "2025-10-26T03:00:00Z",
+ ]
+
+
+class TestCalendarEveryHourCronDstGap(CalendarEveryHourCronDstBase):
+ """Spring-forward (2026-03-29 02:00 CET -> 03:00 CEST): the skipped hour
is not double-counted."""
+
+ DAG_NAME = "test_dag_every_hour_dst_gap"
+ START_DATE = datetime(2026, 3, 1)
+ LAST_RUN_UTC = pendulum.datetime(2026, 3, 28, 22, 0, 0, tz="UTC")
+ EXPECTED_HOURS = [
+ "2026-03-28T23:00:00Z",
+ "2026-03-29T00:00:00Z",
+ "2026-03-29T01:00:00Z",
+ "2026-03-29T02:00:00Z",
+ "2026-03-29T03:00:00Z",
+ ]
+
+
class TestPartitionedCalendar:
"""Calendar tests for partitioned Dags (AIP-76) which use partition_date
instead of logical_date."""