This is an automated email from the ASF dual-hosted git repository.

ephraimbuddy pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new d6cd4595dbd Fix airflow db clean hang on MySQL when delete fails 
(#66296)
d6cd4595dbd is described below

commit d6cd4595dbd55f53d0e0d99a48515c5f0544a610
Author: deepinsight coder <[email protected]>
AuthorDate: Tue Jul 14 01:53:55 2026 -0700

    Fix airflow db clean hang on MySQL when delete fails (#66296)
    
    * Fix airflow db clean hang on MySQL when delete fails (#66177)
    
    * Rename newsfragment to match PR number
    
    The check-newsfragment-pr-number CI job requires the newsfragment
    filename to match the PR number, not the linked issue number.
    
    * Remove newsfragment per review
    
    * Suppress rollback errors so the original delete error propagates
    
    * Explain why the archive table drop reuses the session connection
    
    * Force the FK-failure path in backend regression tests via config copy
    
    Keeps the hang-regression tests meaningful once dag_version cleanup learns
    to skip FK-pinned rows (apache/airflow#68339).
    
    * Guard archive-table cleanup so it can't mask the original delete error
    
    When _do_delete unwinds from a failed DELETE, dropping the archive table in
    the finally block could itself raise (e.g. a dead connection on MySQL). 
Python
    makes a finally-raised error the top-level exception, so that cleanup error
    would replace the original delete failure -- and an OperationalError from 
the
    drop could then be downgraded to a warning by _suppress_with_logging, 
silently
    masking the real IntegrityError. Guard the cleanup so it only propagates on 
the
    success path; on the failure path it is logged and the original error 
survives.
    
    * Update airflow-core/tests/unit/utils/test_db_cleanup.py
    
    Co-authored-by: Ephraim Anierobi <[email protected]>
    
    * Retry CI after provider compat startup failure
    
    The previous run failed because the Mongo provider test fixture could not 
connect to its temporary mongod container. No source change is needed; this 
commit retriggers CI so the PR can use a fresh signal.
    
    ---------
    
    Co-authored-by: Ephraim Anierobi <[email protected]>
    Co-authored-by: Ramachandra Nalam <[email protected]>
---
 airflow-core/src/airflow/utils/db_cleanup.py     |  44 +++-
 airflow-core/tests/unit/utils/test_db_cleanup.py | 318 ++++++++++++++++++++++-
 2 files changed, 353 insertions(+), 9 deletions(-)

diff --git a/airflow-core/src/airflow/utils/db_cleanup.py 
b/airflow-core/src/airflow/utils/db_cleanup.py
index 8197d56c42d..6deb4e78f68 100644
--- a/airflow-core/src/airflow/utils/db_cleanup.py
+++ b/airflow-core/src/airflow/utils/db_cleanup.py
@@ -27,7 +27,7 @@ import csv
 import logging
 import os
 from collections.abc import Generator
-from contextlib import contextmanager
+from contextlib import contextmanager, suppress
 from dataclasses import dataclass
 from types import SimpleNamespace
 from typing import TYPE_CHECKING, Any
@@ -286,6 +286,10 @@ def _do_delete(
         target_table_name = 
f"{ARCHIVE_TABLE_PREFIX}{orm_model.name}__{timestamp_str}{suffix}"
         print(f"Moving data to table {target_table_name}")
         target_table = None
+        # Lets the ``finally`` cleanup below tell the failure path (don't let a
+        # cleanup error mask the original) from the success path (a cleanup 
error
+        # is a real problem and must propagate).
+        error_raised = False
 
         try:
             if dialect_name == "mysql":
@@ -323,13 +327,41 @@ def _do_delete(
             session.execute(delete)
             session.commit()
 
-        except BaseException as e:
-            raise e
+        except BaseException:
+            error_raised = True
+            # Roll back the failed transaction so its locks are released before
+            # the archive table is dropped in the ``finally`` block below.
+            # ``rollback()`` itself can raise (e.g. the connection died); 
suppress
+            # it so it does not shadow the original error being re-raised.
+            with suppress(Exception):
+                session.rollback()
+            raise
         finally:
             if target_table is not None and skip_archive:
-                bind = session.get_bind()
-                target_table.drop(bind=bind)
-                session.commit()
+                # Drop the archive table on the session's own connection. 
Binding
+                # the drop to ``session.get_bind()`` (the Engine) would check 
out a
+                # *second* pooled connection, and on MySQL its ``DROP TABLE`` 
blocks
+                # indefinitely on the metadata lock still held by this 
session's
+                # open transaction when the DELETE above failed -- the ``db 
clean``
+                # hang reported in #66177.
+                try:
+                    target_table.drop(bind=session.connection())
+                    session.commit()
+                except Exception:
+                    # If we are already unwinding from a delete failure, a 
cleanup
+                    # error here must not replace the original exception 
(Python
+                    # makes a ``finally``-raised error the top-level one). Log 
and
+                    # let the original delete error keep propagating. On the 
success
+                    # path (no delete error), a drop/commit failure is a real
+                    # problem, so re-raise it.
+                    if not error_raised:
+                        raise
+                    logger.warning(
+                        "Failed to drop archive table %s while cleaning up 
after a "
+                        "delete failure; propagating the original delete error 
instead.",
+                        target_table_name,
+                        exc_info=True,
+                    )
 
     print("Finished Performing Delete")
 
diff --git a/airflow-core/tests/unit/utils/test_db_cleanup.py 
b/airflow-core/tests/unit/utils/test_db_cleanup.py
index 4d66a96b1a3..8271d38968b 100644
--- a/airflow-core/tests/unit/utils/test_db_cleanup.py
+++ b/airflow-core/tests/unit/utils/test_db_cleanup.py
@@ -17,18 +17,21 @@
 # under the License.
 from __future__ import annotations
 
+import threading
+import time
 from contextlib import suppress
 from importlib import import_module
 from io import StringIO
 from pathlib import Path
-from unittest.mock import MagicMock, mock_open, patch
+from unittest.mock import MagicMock, call, mock_open, patch
 from uuid import uuid4
 
 import pendulum
 import pytest
-from sqlalchemy import func, inspect, select, text
-from sqlalchemy.exc import OperationalError, SQLAlchemyError
+from sqlalchemy import Column, Integer, MetaData, Table, func, inspect, 
literal, select, text
+from sqlalchemy.exc import IntegrityError, OperationalError, SQLAlchemyError
 from sqlalchemy.ext.declarative import DeclarativeMeta
+from sqlalchemy.orm import Session
 
 from airflow import DAG
 from airflow._shared.timezones import timezone
@@ -46,6 +49,7 @@ from airflow.utils.db_cleanup import (
     _build_query,
     _cleanup_table,
     _confirm_drop_archives,
+    _do_delete,
     _dump_table_to_file,
     _get_archived_table_names,
     _TableConfig,
@@ -558,6 +562,167 @@ class TestDBCleanup:
                 # "id" intentionally omitted from extra_columns
             )
 
+    def test_do_delete_rolls_back_before_drop_on_failure(self):
+        session = MagicMock(spec=Session)
+        session.get_bind.return_value.dialect.name = "mysql"
+        session.connection.return_value = object()
+        session.scalars.return_value.one.side_effect = [1, 0]
+        tracker = MagicMock()
+        tracker.attach_mock(session, "session")
+
+        metadata, source_table, target_table, query = 
_build_do_delete_test_objects()
+        delete_failure = IntegrityError("DELETE FROM dag_version", {}, 
Exception("fk violation"))
+        session.execute.side_effect = [None, None, delete_failure]
+
+        with (
+            patch("airflow.utils.db_cleanup.reflect_tables", 
return_value=metadata),
+            patch("airflow.utils.db_cleanup.timezone.utcnow", 
return_value=_delete_test_timestamp()),
+            patch.object(target_table, "drop") as drop_mock,
+        ):
+            tracker.attach_mock(drop_mock, "drop")
+            with pytest.raises(IntegrityError) as exc_info:
+                _do_delete(
+                    query=query,
+                    orm_model=source_table,
+                    skip_archive=True,
+                    session=session,
+                    batch_size=None,
+                )
+
+        assert exc_info.value is delete_failure
+        session.rollback.assert_called_once_with()
+        session.connection.assert_called_once_with()
+        assert session.get_bind.call_count == 1
+        drop_mock.assert_called_once_with(bind=session.connection.return_value)
+
+        rollback_call_index = tracker.mock_calls.index(call.session.rollback())
+        drop_call_index = 
tracker.mock_calls.index(call.drop(bind=session.connection.return_value))
+        commit_call_index = tracker.mock_calls.index(call.session.commit(), 
drop_call_index)
+        assert rollback_call_index < drop_call_index < commit_call_index
+
+    def test_do_delete_propagates_original_error_when_rollback_fails(self):
+        session = MagicMock(spec=Session)
+        session.get_bind.return_value.dialect.name = "mysql"
+        session.connection.return_value = object()
+        session.scalars.return_value.one.side_effect = [1, 0]
+
+        metadata, source_table, target_table, query = 
_build_do_delete_test_objects()
+        delete_failure = IntegrityError("DELETE FROM dag_version", {}, 
Exception("fk violation"))
+        session.execute.side_effect = [None, None, delete_failure]
+        session.rollback.side_effect = OperationalError("ROLLBACK", {}, 
Exception("connection lost"))
+
+        with (
+            patch("airflow.utils.db_cleanup.reflect_tables", 
return_value=metadata),
+            patch("airflow.utils.db_cleanup.timezone.utcnow", 
return_value=_delete_test_timestamp()),
+            patch.object(target_table, "drop") as drop_mock,
+        ):
+            with pytest.raises(IntegrityError) as exc_info:
+                _do_delete(
+                    query=query,
+                    orm_model=source_table,
+                    skip_archive=True,
+                    session=session,
+                    batch_size=None,
+                )
+
+        assert exc_info.value is delete_failure
+        session.rollback.assert_called_once_with()
+        drop_mock.assert_called_once_with(bind=session.connection.return_value)
+
+    @pytest.mark.parametrize(
+        ("skip_archive", "expected_commit_count"),
+        [pytest.param(True, 3, id="skip_archive"), pytest.param(False, 2, 
id="keep_archive")],
+    )
+    def test_do_delete_success_does_not_call_rollback(self, skip_archive, 
expected_commit_count):
+        session = MagicMock(spec=Session)
+        session.get_bind.return_value.dialect.name = "mysql"
+        session.connection.return_value = object()
+        session.scalars.return_value.one.side_effect = [1, 0]
+        session.execute.side_effect = [None, None, None]
+
+        metadata, source_table, target_table, query = 
_build_do_delete_test_objects()
+
+        with (
+            patch("airflow.utils.db_cleanup.reflect_tables", 
return_value=metadata),
+            patch("airflow.utils.db_cleanup.timezone.utcnow", 
return_value=_delete_test_timestamp()),
+            patch.object(target_table, "drop") as drop_mock,
+        ):
+            _do_delete(
+                query=query,
+                orm_model=source_table,
+                skip_archive=skip_archive,
+                session=session,
+                batch_size=None,
+            )
+
+        session.rollback.assert_not_called()
+        assert session.commit.call_count == expected_commit_count
+        if skip_archive:
+            session.connection.assert_called_once_with()
+            
drop_mock.assert_called_once_with(bind=session.connection.return_value)
+        else:
+            session.connection.assert_not_called()
+            drop_mock.assert_not_called()
+
+    def test_do_delete_original_error_survives_archive_drop_failure(self):
+        """On the failure path, a drop/commit error in the finally block must 
not
+        replace the original delete error (nailo2c review, #66296)."""
+        session = MagicMock(spec=Session)
+        session.get_bind.return_value.dialect.name = "mysql"
+        session.connection.return_value = object()
+        session.scalars.return_value.one.side_effect = [1, 0]
+
+        metadata, source_table, target_table, query = 
_build_do_delete_test_objects()
+        delete_failure = IntegrityError("DELETE FROM dag_version", {}, 
Exception("fk violation"))
+        session.execute.side_effect = [None, None, delete_failure]
+        drop_failure = OperationalError("DROP TABLE", {}, Exception("server 
has gone away"))
+
+        with (
+            patch("airflow.utils.db_cleanup.reflect_tables", 
return_value=metadata),
+            patch("airflow.utils.db_cleanup.timezone.utcnow", 
return_value=_delete_test_timestamp()),
+            patch.object(target_table, "drop", side_effect=drop_failure) as 
drop_mock,
+        ):
+            with pytest.raises(IntegrityError) as exc_info:
+                _do_delete(
+                    query=query,
+                    orm_model=source_table,
+                    skip_archive=True,
+                    session=session,
+                    batch_size=None,
+                )
+
+        assert exc_info.value is delete_failure
+        drop_mock.assert_called_once_with(bind=session.connection.return_value)
+
+    def test_do_delete_success_propagates_archive_drop_error(self):
+        """On the success path, a drop/commit failure is a real error and must
+        still surface (the failure-path guard must not swallow it)."""
+        session = MagicMock(spec=Session)
+        session.get_bind.return_value.dialect.name = "mysql"
+        session.connection.return_value = object()
+        session.scalars.return_value.one.side_effect = [1, 0]
+        session.execute.side_effect = [None, None, None]
+
+        metadata, source_table, target_table, query = 
_build_do_delete_test_objects()
+        drop_failure = OperationalError("DROP TABLE", {}, Exception("disk 
full"))
+
+        with (
+            patch("airflow.utils.db_cleanup.reflect_tables", 
return_value=metadata),
+            patch("airflow.utils.db_cleanup.timezone.utcnow", 
return_value=_delete_test_timestamp()),
+            patch.object(target_table, "drop", side_effect=drop_failure),
+        ):
+            with pytest.raises(OperationalError) as exc_info:
+                _do_delete(
+                    query=query,
+                    orm_model=source_table,
+                    skip_archive=True,
+                    session=session,
+                    batch_size=None,
+                )
+
+        assert exc_info.value is drop_failure
+        session.rollback.assert_not_called()
+
     @patch("airflow.utils.db.reflect_tables")
     def test_skip_archive_failure_will_remove_table(self, reflect_tables_mock):
         """
@@ -591,6 +756,67 @@ class TestDBCleanup:
         archived_table_names = _get_archived_table_names(["dag_run"], session)
         assert len(archived_table_names) == 0
 
+    @pytest.mark.backend("mysql")
+    def test_db_clean_does_not_deadlock_on_fk_violation_mysql(self):
+        clean_before_date = _create_pinned_dag_version_cleanup_data(
+            base_date=pendulum.DateTime(2022, 1, 1, 
tzinfo=pendulum.timezone("UTC"))
+        )
+        timeout_seconds = 15
+        finished = threading.Event()
+        result: dict[str, BaseException] = {}
+
+        def cleanup_in_thread():
+            try:
+                with create_session() as session:
+                    _cleanup_table(
+                        **_dag_version_config_without_row_exclusion(),
+                        clean_before_timestamp=clean_before_date,
+                        dry_run=False,
+                        session=session,
+                        table_names=["dag_version"],
+                        skip_archive=True,
+                    )
+            except BaseException as exc:
+                result["exception"] = exc
+            finally:
+                finished.set()
+
+        worker = threading.Thread(target=cleanup_in_thread, daemon=True)
+        start_time = time.monotonic()
+        worker.start()
+
+        assert finished.wait(timeout_seconds), (
+            "dag_version cleanup timed out after "
+            f"{time.monotonic() - start_time:.2f} seconds waiting for the FK 
violation to surface"
+        )
+
+        worker.join(timeout=1)
+        assert not worker.is_alive()
+        assert isinstance(result.get("exception"), IntegrityError)
+
+        with create_session() as session:
+            assert _get_dag_version_archive_table_names(session=session) == []
+
+    @pytest.mark.backend("postgres")
+    def test_db_clean_failure_path_does_not_break_postgres(self):
+        clean_before_date = _create_pinned_dag_version_cleanup_data(
+            base_date=pendulum.DateTime(2022, 1, 1, 
tzinfo=pendulum.timezone("UTC"))
+        )
+
+        with pytest.raises(IntegrityError):
+            with create_session() as session:
+                _cleanup_table(
+                    **_dag_version_config_without_row_exclusion(),
+                    clean_before_timestamp=clean_before_date,
+                    dry_run=False,
+                    session=session,
+                    table_names=["dag_version"],
+                    skip_archive=True,
+                )
+
+        with create_session() as session:
+            assert _get_dag_version_archive_table_names(session=session) == []
+
     def test_no_models_missing(self):
         """
         1. Verify that for all tables in `airflow.models`, we either have them 
enabled in db cleanup,
@@ -1115,3 +1341,89 @@ class TestConnectionTestRequestCleanup:
             assert seeded[state] in survivors, f"{state} row should NOT be 
cleaned up"
         for state in ("success", "failed"):
             assert seeded[state] not in survivors, f"{state} row should be 
cleaned up"
+
+
+def _delete_test_timestamp():
+    return pendulum.DateTime(2024, 1, 2, 3, 4, 5, 
tzinfo=pendulum.timezone("UTC"))
+
+
+def _build_do_delete_test_objects():
+    metadata = MetaData()
+    source_table = Table("dag_version", metadata, Column("id", Integer, 
primary_key=True))
+    target_table = Table(
+        f"{ARCHIVE_TABLE_PREFIX}{source_table.name}__20240102030405",
+        metadata,
+        Column("id", Integer, primary_key=True),
+    )
+    query = select(literal(1).label("id"))
+    return metadata, source_table, target_table, query
+
+
+def _create_pinned_dag_version_cleanup_data(*, base_date):
+    bundle_name = f"testing-{uuid4()}"
+    dag_id = f"test-dag_{uuid4()}"
+
+    with create_session() as session:
+        session.add(DagBundleModel(name=bundle_name))
+        session.flush()
+
+        dag = DAG(dag_id=dag_id)
+        session.add(DagModel(dag_id=dag_id, bundle_name=bundle_name))
+
+        old_dag_version = DagVersion(
+            dag_id=dag_id,
+            version_number=1,
+            bundle_name=bundle_name,
+            created_at=base_date,
+            last_updated=base_date,
+        )
+        new_dag_version = DagVersion(
+            dag_id=dag_id,
+            version_number=2,
+            bundle_name=bundle_name,
+            created_at=base_date.add(minutes=1),
+            last_updated=base_date.add(minutes=1),
+        )
+        session.add_all([old_dag_version, new_dag_version])
+        session.flush()
+
+        dag_run = DagRun(
+            dag_id,
+            run_id="run-1",
+            run_type=DagRunType.MANUAL,
+            start_date=base_date,
+        )
+        task = PythonOperator(task_id="dummy-task", dag=dag, 
python_callable=print)
+        ti = create_task_instance(task, run_id=dag_run.run_id, 
dag_version_id=old_dag_version.id)
+        ti.dag_id = dag_id
+        ti.start_date = base_date
+
+        session.add_all([dag_run, ti])
+        session.commit()
+
+    return base_date.add(days=1)
+
+
+def _get_dag_version_archive_table_names(*, session):
+    return [
+        table_name
+        for table_name in _get_archived_table_names(["dag_version"], session)
+        if table_name.startswith(f"{ARCHIVE_TABLE_PREFIX}dag_version__")
+    ]
+
+
+def _dag_version_config_without_row_exclusion():
+    """Live dag_version cleanup config with row-exclusion filters disabled.
+
+    The backend regression tests above pin ``_do_delete``'s 
rollback-before-drop
+    behaviour, which requires the
+    DELETE to actually hit the ``task_instance.dag_version_id`` FK violation.
+    The live config may exclude FK-pinned rows from the deletion query (see
+    PR #68339), which would turn these tests into no-ops -- so strip any such
+    filters from a copy of the config.
+    """
+    config = dict(config_dict["dag_version"].__dict__)
+    for key in ("extra_filters", "skip_if_referenced"):
+        if key in config:
+            config[key] = None
+    return config

Reply via email to