Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions airflow-core/newsfragments/72370.bugfix.rst
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Fixed failure/retry emails inheriting a stale pinned bundle version for Dag runs with ``disable_bundle_versioning`` enabled, instead of following the unpinned ``dag_run.bundle_version`` like the equivalent task callbacks do.
6 changes: 5 additions & 1 deletion airflow-core/src/airflow/jobs/scheduler_job_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -1680,8 +1680,12 @@ def process_executor_events(
_email_bundle_name = (
ti.dag_version.bundle_name if ti.dag_version else ti.dag_model.bundle_name
)
# Mirror dag_run pinning: if the run wasn't pinned (e.g. dag.disable_bundle_versioning=True),
# leave the callback unpinned so it runs against the same code as the task.
_email_bundle_version = (
ti.dag_version.bundle_version if ti.dag_version else ti.dag_run.bundle_version
ti.dag_version.bundle_version
if ti.dag_version and ti.dag_run.bundle_version is not None
else ti.dag_run.bundle_version
)
_email_version_data = _resolve_version_data(ti.dag_version, ti.dag_run.bundle_version)
# Backfill dag_version_id for legacy tasks (Pydantic requires uuid.UUID).
Expand Down
45 changes: 45 additions & 0 deletions airflow-core/tests/unit/jobs/test_scheduler_job.py
Original file line number Diff line number Diff line change
Expand Up @@ -9445,6 +9445,51 @@ def test_external_kill_callback_bundle_version_follows_dag_run(
assert isinstance(request, TaskCallbackRequest)
assert request.bundle_version == expected_bv

@pytest.mark.parametrize(
("dag_run_bv", "dag_version_bv", "expected_bv"),
[
pytest.param(None, "abc123-sha", None, id="disable_bundle_versioning"),
pytest.param("abc123-sha", "abc123-sha", "abc123-sha", id="versioning_enabled"),
],
)
def test_external_kill_email_bundle_version_follows_dag_run(
self, dag_maker, session, dag_run_bv, dag_version_bv, expected_bv
):
"""
EmailRequest.bundle_version must mirror dag_run.bundle_version, not
dag_version.bundle_version -- the same invariant asserted for
TaskCallbackRequest above by test_external_kill_callback_bundle_version_follows_dag_run.
With disable_bundle_versioning=True the trigger path leaves dag_run.bundle_version=None
even though DagVersion was written with a SHA -- the failure/retry email must inherit
None so the DAG Processor renders it against the same on-disk code as the task, instead
of pinning to a version the run was never pinned to.
"""
with dag_maker(dag_id=f"ext_kill_email_bv_{dag_run_bv or 'none'}", fileloc="/test_path1/"):
EmptyOperator(task_id="t1", email="test@example.com", email_on_failure=True)
dr = dag_maker.create_dagrun(state=DagRunState.RUNNING)

ti = dr.get_task_instance(task_id="t1", session=session)
dag_version = ti.dag_version
dag_version.bundle_version = dag_version_bv
dr.bundle_version = dag_run_bv
ti.state = State.QUEUED
session.merge(dag_version)
session.merge(dr)
session.merge(ti)
session.commit()

executor = MockExecutor(do_update=False)
scheduler_job = Job()
self.job_runner = SchedulerJobRunner(scheduler_job, executors=[executor])
executor.event_buffer[ti.key] = State.FAILED, None

self.job_runner._process_executor_events(executor=executor, session=session)

self.job_runner.executor.callback_sink.send.assert_called_once()
request = self.job_runner.executor.callback_sink.send.call_args[0][0]
assert isinstance(request, EmailRequest)
assert request.bundle_version == expected_bv

def test_heartbeat_timeout_callback_bundle_version_follows_dag_run(self, dag_maker, session):
"""
Same invariant as the external-kill path, exercised through
Expand Down