diff --git a/airflow-core/newsfragments/72370.bugfix.rst b/airflow-core/newsfragments/72370.bugfix.rst new file mode 100644 index 0000000000000..3c6515e8215c1 --- /dev/null +++ b/airflow-core/newsfragments/72370.bugfix.rst @@ -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. diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index 127f385f0e014..bbed2c2905de7 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -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). diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py b/airflow-core/tests/unit/jobs/test_scheduler_job.py index a7b332e763f1c..55b7f45cdecc7 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -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