ashb commented on code in PR #73221:
URL: https://github.com/apache/airflow/pull/73221#discussion_r4102790002
##########
airflow-core/src/airflow/dag_processing/processor.py:
##########
@@ -543,6 +543,10 @@ def _execute_email_callbacks(dagbag: DagBag, request:
EmailRequest, log: Filteri
task=task,
_ti_context_from_server=ctx_from_server,
max_tries=ctx_from_server.max_tries,
+ # The callback request carries no state; the email type is the reason
it fired.
+ state=(
+ TaskInstanceState.FAILED if request.email_type == "failure" else
TaskInstanceState.UP_FOR_RETRY
+ ),
Review Comment:
We already have `request.email_type == "failure"` conditions on L524.
Instead of an inline condition, I think we shold set a `state` local var in
those existing conditions.
##########
providers/smtp/tests/unit/smtp/notifications/test_smtp.py:
##########
@@ -163,14 +168,103 @@ def test_notifier_with_defaults(self,
mock_smtphook_hook, create_dag_without_db,
mock_smtphook_hook.return_value.__enter__().send_email_smtp.assert_called_once_with(
from_email=TEST_SENDER,
to=TEST_RECEIVER,
- subject=f"DAG {TEST_DAG_ID} - Task {TEST_TASK_ID} - Run ID
{TEST_RUN_ID} in State {TEST_TASK_STATE}",
+ subject=f"[Airflow] {TEST_DAG_ID}.{TEST_TASK_ID}
{TEST_TASK_STATE.value} - Run {TEST_RUN_ID}",
html_content=mock.ANY,
smtp_conn_id=SMTP_CONN_ID,
**DEFAULT_EMAIL_PARAMS,
)
content =
mock_smtphook_hook.return_value.__enter__().send_email_smtp.call_args.kwargs["html_content"]
assert f"{TRY_NUMBER} of 1" in content
+ @pytest.mark.parametrize(
+ ("state", "expected_state", "expected_banner"),
+ [
+ pytest.param(TaskInstanceState.FAILED, "failed", "#dc2626",
id="failed"),
+ pytest.param(TaskInstanceState.SUCCESS, "success", "#334155",
id="success"),
+ pytest.param(None, "unknown", "#334155", id="no-state"),
+ ],
+ )
+ @mock.patch("airflow.providers.smtp.notifications.smtp.SmtpHook")
Review Comment:
All mocks that get added/changed should have spec/autospec please.
##########
providers/smtp/tests/unit/smtp/notifications/test_smtp.py:
##########
@@ -163,14 +168,103 @@ def test_notifier_with_defaults(self,
mock_smtphook_hook, create_dag_without_db,
mock_smtphook_hook.return_value.__enter__().send_email_smtp.assert_called_once_with(
from_email=TEST_SENDER,
to=TEST_RECEIVER,
- subject=f"DAG {TEST_DAG_ID} - Task {TEST_TASK_ID} - Run ID
{TEST_RUN_ID} in State {TEST_TASK_STATE}",
+ subject=f"[Airflow] {TEST_DAG_ID}.{TEST_TASK_ID}
{TEST_TASK_STATE.value} - Run {TEST_RUN_ID}",
Review Comment:
This `[Airflow]` prefix change seems subjective, and not key to the core
change. Please undo that part (and introduce it as a separate change if you
think it's worth it)
##########
task-sdk/src/airflow/sdk/execution_time/task_runner.py:
##########
@@ -2152,6 +2168,8 @@ def _send_error_email_notification(
"exception_html": exception_html,
"try_number": ti.try_number,
"max_tries": ti.max_tries,
+ # Pass the value, not the enum: str() on it renders as
"TaskInstanceState.FAILED".
+ "task_state": ti.state.value if ti.state else "unknown",
Review Comment:
`ti.state` of None should not show "unknown", but likey "None" which is what
the UI does.
##########
task-sdk/tests/task_sdk/execution_time/test_task_runner.py:
##########
@@ -4183,9 +4205,13 @@ def execute(self, context):
kwargs = mock_smtp_notifier.call_args.kwargs
assert kwargs["from_email"] == self.FROM
assert kwargs["to"] == emails
+ assert (
+ kwargs["subject"]
+ == "[Airflow] {{ti.dag_id}}.{{ti.task_id}}
{{task_state}} - Run {{ti.run_id}}"
+ )
assert (
kwargs["html_content"]
- == 'Try {{try_number}} out of {{max_tries +
1}}<br>Exception:<br>{{exception_html}}<br>Log: <a
href="{{ti.log_url}}">Link</a><br>Host: {{ti.hostname}}<br>Mark success: <a
href="{{ti.mark_success_url}}">Link</a><br>'
+ == 'Dag: {{ti.dag_id}}<br>Task: {{ti.task_id}}<br>Run:
{{ti.run_id}}<br>State: {{task_state}}<br>Try: {{try_number}} out of
{{max_tries + 1}}<br>{% if ti.start_date is defined and ti.start_date
%}Started: {{ti.start_date}}<br>{% endif %}{% if ti.end_date is defined and
ti.end_date %}Ended: {{ti.end_date}}<br>{% endif
%}Exception:<br>{{exception_html}}<br>Log: <a
href="{{ti.log_url}}">Link</a><br>Host: {{ti.hostname}}<br>'
)
Review Comment:
This test feels low value tbh, it doesn't assert anything other than "the
test has been kept up to date with the code".
`assert "html_content" in kwargs` is probably sufficient
##########
airflow-core/newsfragments/73221.bugfix.rst:
##########
Review Comment:
Bug fix newsfragments should be much, much shorter. Please hand write this
instead of letting an LLM write it.
##########
task-sdk/src/airflow/sdk/execution_time/task_runner.py:
##########
@@ -2124,7 +2135,7 @@ def _send_error_email_notification(
subject = Path(subject_template_file).read_text()
else:
# Fallback to default
- subject = "Airflow alert: {{ti}}"
+ subject = "[Airflow] {{ti.dag_id}}.{{ti.task_id}} {{task_state}} - Run
{{ti.run_id}}"
Review Comment:
Ditto on the prefix.
##########
task-sdk/tests/task_sdk/execution_time/test_task_runner.py:
##########
@@ -4293,6 +4325,42 @@ def execute(self, context):
)
assert kwargs["from_email"] == self.FROM
+ def test_default_email_body_renders_without_dates(self, create_runtime_ti,
mock_supervisor_comms):
+ from airflow.sdk.execution_time.task_runner import
_send_error_email_notification
+
+ task = BaseOperator(task_id="callback_task",
email=["[email protected]"], email_on_failure=True)
+ runtime_ti = create_runtime_ti(task=task)
+ # The Dag-processor callback path builds the task instance from the
callback request alone,
+ # which carries no dates, and the Dag's Jinja environment is
StrictUndefined.
+ callback_ti = RuntimeTaskInstance.model_construct(
+ id=runtime_ti.id,
+ task_id=runtime_ti.task_id,
+ dag_id=runtime_ti.dag_id,
+ run_id=runtime_ti.run_id,
+ try_number=runtime_ti.try_number,
+ dag_version_id=runtime_ti.dag_version_id,
+ task=task,
+ _ti_context_from_server=None,
+ max_tries=0,
+ )
+ context = callback_ti.get_template_context()
+ log = mock.MagicMock()
+
+ with conf_vars({("email", "from_email"): self.FROM}):
+ with
mock.patch("airflow.providers.smtp.notifications.smtp.SmtpNotifier") as
mock_smtp_notifier:
+ _send_error_email_notification(task, callback_ti, context,
ValueError("boom"), log)
+
+ kwargs = mock_smtp_notifier.call_args.kwargs
+ email_context = mock_smtp_notifier.return_value.call_args.args[0]
+ env = task.dag.get_template_env()
+
+ assert env.from_string(kwargs["subject"]).render(email_context) == (
+ f"[Airflow] {callback_ti.dag_id}.{callback_ti.task_id} unknown -
Run {callback_ti.run_id}"
Review Comment:
I'd rather we provided specific values to `create_runtime_ti` so that this
test can be something like
```suggestion
f"[Airflow] test_dag.example_task unknown - Run
manual_20260925T..."
```
etc. It feels like relying on f-strings here could hide any manner of breaks
without noticing.
--
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]