From 69649b1dacb002591633fd879f1566b78d77f3fa Mon Sep 17 00:00:00 2001 From: Sean Ghaeli Date: Tue, 23 Jun 2026 19:08:11 +0000 Subject: [PATCH] Resolve VariableInterval deadlines safely at DagRun creation A ``DeadlineAlert`` configured with a ``VariableInterval`` is resolved when the scheduler creates a DagRun, inside ``DAG._process_dagrun_deadline_alerts`` (which runs under the ``prohibit_commit`` guard). Two problems are fixed: 1. Resolution now goes through the full secrets chain (env vars, configured secrets backends, then the metadata DB) via a dedicated ``_resolve_variable_interval`` helper, rather than ``Variable.get`` / ``begin_nested``. ``Variable.get`` and a SAVEPOINT release both commit on the scheduler's session, tripping ``prohibit_commit`` ("UNEXPECTED COMMIT") and silently dropping deadlines for every scheduled DagRun. The helper passes the scheduler session through to the metastore backend so the DB read happens without committing, and reading via the secrets chain (not the variable table directly) means ``AIRFLOW_VAR_*`` env vars and secrets backends resolve too. 2. Each deadline alert is isolated with a plain ``try``/``except`` (NOT ``begin_nested``). Creating a deadline is auxiliary to creating the DagRun; a single bad alert -- a missing/invalid backing Variable, or an undecodable serialized blob -- must never abort the DagRun and stop the DAG scheduling. ``VariableInterval.resolve`` is split into ``resolve`` + ``coerce_to_timedelta`` so the scheduler-side reader reuses the exact same validation (including the OverflowError -> ValueError translation) without going through ``Variable.get``. Generated-by: Claude Code (Sonnet/Opus via Claude Code) on behalf of Sean Ghaeli --- airflow-core/newsfragments/68917.bugfix.rst | 1 + .../airflow/serialization/definitions/dag.py | 126 +++++++++++------ airflow-core/tests/unit/models/test_dagrun.py | 133 +++++++++++------- .../metrics/metrics_template.yaml | 6 + .../src/airflow/sdk/definitions/deadline.py | 21 ++- .../task_sdk/definitions/test_deadline.py | 15 ++ 6 files changed, 212 insertions(+), 90 deletions(-) create mode 100644 airflow-core/newsfragments/68917.bugfix.rst diff --git a/airflow-core/newsfragments/68917.bugfix.rst b/airflow-core/newsfragments/68917.bugfix.rst new file mode 100644 index 0000000000000..71cbdf9b05ef1 --- /dev/null +++ b/airflow-core/newsfragments/68917.bugfix.rst @@ -0,0 +1 @@ +Deadline alerts using a ``VariableInterval`` no longer risk aborting DagRun creation. The interval is now resolved through the full secrets chain (environment variables, configured secrets backends, then the metadata database) on the scheduler's own session, so ``AIRFLOW_VAR_*`` and secrets-backend-backed Variables resolve correctly and the read does not commit inside the scheduler's ``prohibit_commit`` guard. Each deadline alert is also isolated: a single unresolvable or undecodable alert is logged and skipped instead of preventing the DagRun from being created. diff --git a/airflow-core/src/airflow/serialization/definitions/dag.py b/airflow-core/src/airflow/serialization/definitions/dag.py index 8ed0fee2ccabd..0eef340273901 100644 --- a/airflow-core/src/airflow/serialization/definitions/dag.py +++ b/airflow-core/src/airflow/serialization/definitions/dag.py @@ -31,8 +31,9 @@ from sqlalchemy import func, or_, select, tuple_ from airflow._shared.observability.metrics import stats +from airflow._shared.secrets_backend.base import call_secrets_backend_method from airflow._shared.timezones.timezone import coerce_datetime -from airflow.configuration import conf as airflow_conf +from airflow.configuration import conf as airflow_conf, ensure_secrets_loaded from airflow.exceptions import ( AirflowException, DagNotPartitionedError, @@ -51,6 +52,7 @@ from airflow.models.taskinstancekey import TaskInstanceKey from airflow.models.tasklog import LogTemplate from airflow.sdk.definitions.deadline import VariableInterval +from airflow.secrets.metastore import MetastoreBackend from airflow.serialization.decoders import decode_deadline_alert from airflow.serialization.definitions.deadline import DeadlineAlertFields, SerializedReferenceModels from airflow.serialization.definitions.param import SerializedParamsDict @@ -741,51 +743,93 @@ def _process_dagrun_deadline_alerts( if not deadline_alert: continue - deserialized_deadline_alert = decode_deadline_alert( - { - Encoding.TYPE: DAT.DEADLINE_ALERT, - Encoding.VAR: { - DeadlineAlertFields.REFERENCE: deadline_alert.reference, - DeadlineAlertFields.INTERVAL: deadline_alert.interval, - DeadlineAlertFields.CALLBACK: deadline_alert.callback_def, - }, - } - ) + # Deadline creation is best-effort. A failure here must not prevent the DagRun + # itself from being created. Use a plain try/except rather than + # ``session.begin_nested()`` since ``create_dagrun`` runs under + # ``prohibit_commit`` and releasing a SAVEPOINT would trip that guard. + try: + deserialized_deadline_alert = decode_deadline_alert( + { + Encoding.TYPE: DAT.DEADLINE_ALERT, + Encoding.VAR: { + DeadlineAlertFields.REFERENCE: deadline_alert.reference, + DeadlineAlertFields.INTERVAL: deadline_alert.interval, + DeadlineAlertFields.CALLBACK: deadline_alert.callback_def, + }, + } + ) - interval = deserialized_deadline_alert.interval + interval = deserialized_deadline_alert.interval - if isinstance(interval, VariableInterval): - interval = interval.resolve() + if isinstance(interval, VariableInterval): + interval = self._resolve_variable_interval(interval, session=session) - if isinstance(deserialized_deadline_alert.reference, SerializedReferenceModels.TYPES.DAGRUN): - deadline_time = deserialized_deadline_alert.reference.evaluate_with( - session=session, - interval=interval, - # TODO : Pretty sure we can drop these last two; verify after testing is complete - dag_id=self.dag_id, - run_id=orm_dagrun.run_id, - ) + if isinstance(deserialized_deadline_alert.reference, SerializedReferenceModels.TYPES.DAGRUN): + deadline_time = deserialized_deadline_alert.reference.evaluate_with( + session=session, + interval=interval, + # TODO : Pretty sure we can drop these last two; verify after testing is complete + dag_id=self.dag_id, + run_id=orm_dagrun.run_id, + ) - if deadline_time is not None: - session.add( - Deadline( - deadline_time=deadline_time, - callback=deserialized_deadline_alert.callback, - dagrun_id=orm_dagrun.id, - deadline_alert_id=deadline_alert.id, - dag_id=orm_dagrun.dag_id, - bundle_name=orm_dagrun.dag_model.bundle_name, + if deadline_time is not None: + session.add( + Deadline( + deadline_time=deadline_time, + callback=deserialized_deadline_alert.callback, + dagrun_id=orm_dagrun.id, + deadline_alert_id=deadline_alert.id, + dag_id=orm_dagrun.dag_id, + bundle_name=orm_dagrun.dag_model.bundle_name, + ) ) - ) - team_name = ( - DagModel.get_team_name(self.dag_id, session=session) - if airflow_conf.getboolean("core", "multi_team") - else None - ) - stats.incr( - "deadline_alerts.deadline_created", - tags=prune_dict({"dag_id": self.dag_id, "team_name": team_name}), - ) + team_name = ( + DagModel.get_team_name(self.dag_id, session=session) + if airflow_conf.getboolean("core", "multi_team") + else None + ) + stats.incr( + "deadline_alerts.deadline_created", + tags=prune_dict({"dag_id": self.dag_id, "team_name": team_name}), + ) + except Exception: + log.exception( + "Failed to create deadline for alert %s on DagRun %s (dag_id=%s); " + "skipping this deadline, the DagRun is unaffected", + getattr(deadline_alert, "id", ""), + orm_dagrun.run_id, + self.dag_id, + ) + stats.incr("deadline_alerts.deadline_creation_failed", tags={"dag_id": self.dag_id}) + + @staticmethod + def _resolve_variable_interval(interval: VariableInterval, *, session: Session) -> datetime.timedelta: + """ + Resolve a ``VariableInterval`` to a concrete ``timedelta`` at DagRun creation. + + The Variable is resolved using the standard secrets lookup order. The scheduler + session is passed to the metastore backend to avoid creating a new session + during DagRun creation. + + :param interval: The ``VariableInterval`` to resolve. + :param session: Scheduler session used for metadata database lookups. + :return: The resolved ``timedelta``. + :raises ValueError: If the Variable cannot be resolved or converted to a valid ``timedelta``. + """ + for backend in ensure_secrets_loaded(): + value = call_secrets_backend_method( + backend.get_variable, + team_name=None, + key=interval.key, + **({"session": session} if isinstance(backend, MetastoreBackend) else {}), + ) + if value is not None: + return interval.coerce_to_timedelta(value) + raise ValueError( + f"VariableInterval '{interval.key}' could not be resolved from any " + f"secrets backend, environment variable, or the metadata database" + ) @provide_session def set_task_instance_state( diff --git a/airflow-core/tests/unit/models/test_dagrun.py b/airflow-core/tests/unit/models/test_dagrun.py index c46c30d1194e6..457f1c8e5a8d8 100644 --- a/airflow-core/tests/unit/models/test_dagrun.py +++ b/airflow-core/tests/unit/models/test_dagrun.py @@ -76,8 +76,6 @@ ) from airflow.sdk.definitions.callback import AsyncCallback from airflow.sdk.definitions.deadline import DeadlineAlert, DeadlineReference, VariableInterval -from airflow.sdk.definitions.variable import Variable -from airflow.sdk.exceptions import AirflowRuntimeError from airflow.serialization.definitions.deadline import SerializedReferenceModels from airflow.serialization.serialized_objects import LazyDeserializedDAG from airflow.settings import get_policy_plugin_manager @@ -1390,16 +1388,20 @@ def test_dag_run_dag_versions_with_null_created_dag_version(self, dag_maker, ses VariableInterval("my_key"), ], ) - @mock.patch.object(Variable, "get") @mock.patch.object(Deadline, "prune_deadlines") - def test_dagrun_success_deadline(self, _, mock_get, interval, session, deadline_test_dag): + def test_dagrun_success_deadline(self, _, interval, session, deadline_test_dag): def on_success_callable(context): assert context["dag_run"].dag_id == "test_dag" - future_date = datetime.datetime.now() + datetime.timedelta(days=365) + future_date = datetime.datetime(2037, 1, 1, tzinfo=datetime.timezone.utc) - # First value used during resolution - mock_get.return_value = "5" + if isinstance(interval, VariableInterval): + # Seed via the metastore model, not the SDK Variable (whose set() routes through + # SUPERVISOR_COMMS), so the row lands in the variable table on this session. + from airflow.models.variable import Variable as VariableModel + + VariableModel.set(key="my_key", value="5", session=session) + session.flush() scheduler_dag = deadline_test_dag( deadline=DeadlineAlert( @@ -1509,71 +1511,108 @@ def test_dagrun_success_handles_empty_deadline_list(self, mock_prune, dag_maker, mock_prune.assert_not_called() assert dag_run.state == DagRunState.SUCCESS - @mock.patch.object(Variable, "get") @mock.patch.object(Deadline, "prune_deadlines") - def test_dagrun_deadline_variable_interval_stable(self, _, mock_get, session, deadline_test_dag): - future_date = datetime.datetime.now() + datetime.timedelta(days=365) - - # First value used during resolution. - mock_get.return_value = "60" + def test_dagrun_deadline_variable_interval_missing_variable_is_isolated( + self, _, session, deadline_test_dag + ): + """A VariableInterval whose backing Variable is missing must NOT abort DagRun creation. + + ``VariableInterval.resolve()`` raises ``ValueError`` for a missing/invalid Variable. + That resolution happens inside ``_process_dagrun_deadline_alerts`` during + ``create_dagrun``; previously the error propagated out and aborted the whole run, + silently stopping the DAG from scheduling. The per-alert ``try``/``except`` now isolates + the failure: the DagRun is created, the bad deadline is skipped (logged), and no Deadline + row is written. (Isolation must NOT use ``begin_nested`` here -- ``create_dagrun`` runs + under the scheduler ``prohibit_commit`` guard, where a SAVEPOINT release would raise + ``UNEXPECTED COMMIT`` and skip every scheduled DagRun's deadlines.) + """ + future_date = datetime.datetime(2037, 1, 1, tzinfo=datetime.timezone.utc) scheduler_dag = deadline_test_dag( deadline=DeadlineAlert( reference=DeadlineReference.FIXED_DATETIME(future_date), - interval=VariableInterval("my_key"), + interval=VariableInterval("missing_key"), callback=AsyncCallback(empty_callback_for_deadline), ), ) dag_run = self.create_dag_run( dag=scheduler_dag, - task_states={"task_1": TaskInstanceState.SUCCESS, "task_2": TaskInstanceState.SUCCESS}, + task_states={"task_1": TaskInstanceState.SUCCESS}, session=session, ) - dag_run.dag = scheduler_dag - - # First update resolve interval to "5". - dag_run.update_state(session=session) + assert dag_run is not None deadline = session.execute(select(Deadline)).scalars().one_or_none() - first_deadline_time = deadline.deadline_time + assert deadline is None - # Change Variable value after resolution. - mock_get.return_value = "120" + @mock.patch.object(Deadline, "prune_deadlines") + def test_dagrun_deadline_variable_interval_resolves_from_env_var( + self, _, session, deadline_test_dag, monkeypatch + ): + """A VariableInterval backed by an ``AIRFLOW_VAR_*`` env var (no DB row) must resolve. - # Run again (This should not change existing deadline). - dag_run.update_state(session=session) + Regression guard: the scheduler-side resolver must go through the full secrets chain + (env vars + secrets backends + metadata DB), not read only the ``variable`` table. A + table-only read returns None for an env/secrets-backed Variable, and the per-alert + ``except`` then silently drops the deadline. Here the Variable lives ONLY in the + environment, so a correct resolver creates the Deadline and a regressed one drops it. + """ + # Variable lives only as an env var, never in the variable table. Values are seconds. + monkeypatch.setenv("AIRFLOW_VAR_ENV_INTERVAL_KEY", "7") + future_date = datetime.datetime(2037, 1, 1, tzinfo=datetime.timezone.utc) + + scheduler_dag = deadline_test_dag( + deadline=DeadlineAlert( + reference=DeadlineReference.FIXED_DATETIME(future_date), + interval=VariableInterval("env_interval_key"), + callback=AsyncCallback(empty_callback_for_deadline), + ), + ) + + dag_run = self.create_dag_run( + dag=scheduler_dag, + task_states={"task_1": TaskInstanceState.SUCCESS}, + session=session, + ) + assert dag_run is not None deadline = session.execute(select(Deadline)).scalars().one_or_none() - assert deadline.deadline_time == first_deadline_time + assert deadline is not None + assert deadline.deadline_time == future_date + datetime.timedelta(seconds=7) @mock.patch.object(Deadline, "prune_deadlines") - def test_dagrun_deadline_variable_interval_missing_variable_fails(self, _, session, deadline_test_dag): - mock_err = mock.Mock() - mock_err.error.value = "MISSING_DEADLINE" - mock_err.detail = "missing deadline" + def test_dagrun_deadline_decode_failure_is_isolated(self, _, session, deadline_test_dag): + """A deadline alert that fails to *decode* must NOT abort DagRun creation either. + + ``decode_deadline_alert`` can raise for a malformed/legacy serialized blob (e.g. a + None interval after a partial downgrade, an invalid interval type, or a corrupt + reference dict). That decode happens inside the per-alert ``try``/``except``, so the + failure is isolated just like a resolve-time failure: the DagRun is created and the bad + deadline skipped, rather than the corrupt row taking down scheduling for the whole DAG. + """ + future_date = datetime.datetime(2037, 1, 1, tzinfo=datetime.timezone.utc) + scheduler_dag = deadline_test_dag( + deadline=DeadlineAlert( + reference=DeadlineReference.FIXED_DATETIME(future_date), + interval=datetime.timedelta(hours=1), + callback=AsyncCallback(empty_callback_for_deadline), + ), + ) - with mock.patch.object( - Variable, - "get", - side_effect=AirflowRuntimeError(mock_err), + with mock.patch( + "airflow.serialization.definitions.dag.decode_deadline_alert", + side_effect=ValueError("corrupt deadline alert blob"), ): - future_date = datetime.datetime.now() + datetime.timedelta(days=365) - - scheduler_dag = deadline_test_dag( - deadline=DeadlineAlert( - reference=DeadlineReference.FIXED_DATETIME(future_date), - interval=VariableInterval("missing_key"), - callback=AsyncCallback(empty_callback_for_deadline), - ), + dag_run = self.create_dag_run( + dag=scheduler_dag, + task_states={"task_1": TaskInstanceState.SUCCESS}, + session=session, ) - with pytest.raises(ValueError, match="not found"): - self.create_dag_run( - dag=scheduler_dag, - task_states={"task_1": TaskInstanceState.SUCCESS}, - session=session, - ) + assert dag_run is not None + deadline = session.execute(select(Deadline)).scalars().one_or_none() + assert deadline is None @pytest.mark.parametrize( diff --git a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml index a561db2d7aef6..35777f8515754 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -310,6 +310,12 @@ metrics: legacy_name: "-" name_variables: [] + - name: "deadline_alerts.deadline_creation_failed" + description: "Number of deadline alerts that could not be created (interval/reference resolution raised)" + type: "counter" + legacy_name: "-" + name_variables: [] + - name: "deadline_alerts.deadline_missed" description: "Number of deadline alerts that fired because a Dag run missed its deadline" type: "counter" diff --git a/task-sdk/src/airflow/sdk/definitions/deadline.py b/task-sdk/src/airflow/sdk/definitions/deadline.py index a9da3dd3d3ea2..04bd43cb7108d 100644 --- a/task-sdk/src/airflow/sdk/definitions/deadline.py +++ b/task-sdk/src/airflow/sdk/definitions/deadline.py @@ -389,9 +389,19 @@ def resolve(self) -> timedelta: value = Variable.get(self.key) except AirflowRuntimeError as e: raise ValueError(f"VariableInterval '{self.key}' not found") from e + return self.coerce_to_timedelta(value) + def coerce_to_timedelta(self, value: str | int | float | None) -> timedelta: + """ + Validate a raw Variable value and convert it into a ``timedelta``. + + Split out from :meth:`resolve` so callers that already hold the Variable value (e.g. a + scheduler-side reader that must fetch it on its own session — see + ``DAG._process_dagrun_deadline_alerts`` — to avoid committing inside ``prohibit_commit``) + can reuse the exact same validation without going through ``Variable.get``. + """ try: - seconds = int(value) + seconds = int(value) # type: ignore[arg-type] # None/non-numeric handled by the except below except (TypeError, ValueError) as e: raise ValueError( f"VariableInterval '{self.key}' must be an integer (seconds), got: {value!r}" @@ -400,4 +410,11 @@ def resolve(self) -> timedelta: if seconds <= 0: raise ValueError(f"VariableInterval '{self.key}' must be > 0, got: {seconds}") - return timedelta(seconds=seconds) + try: + return timedelta(seconds=seconds) + except OverflowError as e: + # A huge value overflows ``timedelta`` (OverflowError, not a ValueError subclass). + # Translate it to a clean ValueError so callers get a consistent error. + raise ValueError( + f"VariableInterval '{self.key}' is too large to be a valid interval: {seconds} seconds" + ) from e diff --git a/task-sdk/tests/task_sdk/definitions/test_deadline.py b/task-sdk/tests/task_sdk/definitions/test_deadline.py index b104980e4c986..8c24fb11b1524 100644 --- a/task-sdk/tests/task_sdk/definitions/test_deadline.py +++ b/task-sdk/tests/task_sdk/definitions/test_deadline.py @@ -212,3 +212,18 @@ def test_resolve_invalid(self, mocker, value, raise_runtime, match): with pytest.raises(ValueError, match=match): interval.resolve() + + @pytest.mark.parametrize( + ("value", "match"), + [ + ("abc", "must be an integer"), + ("", "must be an integer"), + (None, "must be an integer"), + ("0", "must be > 0"), + ("-5", "must be > 0"), + ("99999999999999", "too large to be a valid interval"), + ], + ) + def test_coerce_to_timedelta_invalid(self, value, match): + with pytest.raises(ValueError, match=match): + VariableInterval(key="k").coerce_to_timedelta(value)