seanghaeli commented on code in PR #68917:
URL: https://github.com/apache/airflow/pull/68917#discussion_r3787932135
##########
airflow-core/tests/unit/models/test_dagrun.py:
##########
@@ -1509,71 +1511,121 @@ 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")
+ @pytest.mark.parametrize(
+ ("interval", "failure"),
+ [
+ pytest.param(VariableInterval("missing_key"), nullcontext(),
id="missing_variable"),
+ pytest.param(
+ datetime.timedelta(hours=1),
+ mock.patch(
+
"airflow.serialization.definitions.dag.decode_deadline_alert",
+ autospec=True,
+ side_effect=ValueError("corrupt deadline alert blob"),
+ ),
+ id="decode_failure",
+ ),
+ pytest.param(
+ datetime.timedelta(hours=1),
+ mock.patch.object(
+ SerializedReferenceModels.FixedDatetimeDeadline,
+ "evaluate_with",
+ autospec=True,
+ side_effect=RuntimeError("evaluate_with failed"),
+ ),
+ id="evaluate_with_failure",
+ ),
+ ],
+ )
@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)
+ def test_dagrun_deadline_failure_is_isolated(self, _, interval, failure,
session, deadline_test_dag):
+ """A failure while creating any single deadline must not abort DagRun
creation."""
+ 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=interval,
+ callback=AsyncCallback(empty_callback_for_deadline),
+ ),
+ )
+
+ with failure:
+ dag_run = self.create_dag_run(
+ dag=scheduler_dag,
+ task_states={"task_1": TaskInstanceState.SUCCESS},
+ session=session,
+ )
- # First value used during resolution.
- mock_get.return_value = "60"
+ assert dag_run is not None
+ assert session.execute(select(Deadline)).scalars().one_or_none() is
None
+
+ @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."""
+ 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("my_key"),
+ 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, "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)
-
- deadline = session.execute(select(Deadline)).scalars().one_or_none()
- first_deadline_time = deadline.deadline_time
-
- # Change Variable value after resolution.
- mock_get.return_value = "120"
-
- # Run again (This should not change existing deadline).
- dag_run.update_state(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)
+ @pytest.mark.parametrize(
+ ("multi_team", "team_name"),
+ [
+ pytest.param("true", "team_alpha", id="team_scoped"),
+ pytest.param("false", None, id="global"),
+ ],
+ )
@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_variable_interval_scoped_to_team(
+ self, _, multi_team, team_name, session, deadline_test_dag
+ ):
+ """A VariableInterval must resolve against the DagRun's team, not the
global scope."""
+ future_date = datetime.datetime(2037, 1, 1,
tzinfo=datetime.timezone.utc)
- with mock.patch.object(
- Variable,
- "get",
- side_effect=AirflowRuntimeError(mock_err),
- ):
- future_date = datetime.datetime.now() +
datetime.timedelta(days=365)
+ scheduler_dag = deadline_test_dag(
+ deadline=DeadlineAlert(
+ reference=DeadlineReference.FIXED_DATETIME(future_date),
+ interval=VariableInterval("team_interval_key"),
+ callback=AsyncCallback(empty_callback_for_deadline),
+ ),
+ )
- scheduler_dag = deadline_test_dag(
- deadline=DeadlineAlert(
- reference=DeadlineReference.FIXED_DATETIME(future_date),
- interval=VariableInterval("missing_key"),
- callback=AsyncCallback(empty_callback_for_deadline),
- ),
+ with (
+ conf_vars({("core", "multi_team"): multi_team}),
+ mock.patch("airflow.models.dag.DagModel.get_team_name",
return_value=team_name),
+ mock.patch(
+
"airflow.serialization.definitions.dag.call_secrets_backend_method",
+ return_value="5",
+ ) as mock_call,
+ ):
+ 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
+ # The team the DagRun belongs to must be forwarded to the backend
lookup.
+ mock_call.assert_called_once()
+ assert mock_call.call_args.kwargs["team_name"] == team_name
Review Comment:
Good catch, definitely an oversight by me. Addressed in the latest commit
--
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]