SameerMesiah97 commented on code in PR #68917:
URL: https://github.com/apache/airflow/pull/68917#discussion_r3695583135
##########
airflow-core/src/airflow/serialization/definitions/dag.py:
##########
@@ -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(
Review Comment:
Have you added a test for a scenario where `evaluate_with` fails or
otherwise raises an error?
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -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.
Review Comment:
The above feedback applies here too.
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -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.
Review Comment:
Only keep the first line. Like I mentioned previously, we do not need long
docstrings. Additional context should be covered in comments.
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -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.
Review Comment:
Here too.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]