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

kaxil 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 ef54cd661c4 Fix deactivation of stale zip-packaged DAGs (#70586)
ef54cd661c4 is described below

commit ef54cd661c4441650ca5294a1042ee85c1570cbe
Author: Sam Wheating <[email protected]>
AuthorDate: Sat Aug 8 11:03:56 2026 +0100

    Fix deactivation of stale zip-packaged DAGs (#70586)
---
 airflow-core/src/airflow/dag_processing/manager.py |  6 +-
 .../tests/unit/dag_processing/test_manager.py      | 69 +++++++++++++++++++++-
 2 files changed, 72 insertions(+), 3 deletions(-)

diff --git a/airflow-core/src/airflow/dag_processing/manager.py 
b/airflow-core/src/airflow/dag_processing/manager.py
index df0f8f87481..a4d4917958c 100644
--- a/airflow-core/src/airflow/dag_processing/manager.py
+++ b/airflow-core/src/airflow/dag_processing/manager.py
@@ -506,7 +506,11 @@ class DagFileProcessorManager(LoggingMixin):
             # When the Dag's last_parsed_time is more than the 
stale_dag_threshold older than the
             # Dag file's last_finish_time, the Dag is considered stale as has 
apparently been removed from the file,
             # This is especially relevant for Dag files that generate Dags in 
a dynamic manner.
-            file_info = DagFileInfo(rel_path=Path(dag.relative_fileloc), 
bundle_name=dag.bundle_name)
+            rel_path = Path(dag.relative_fileloc)
+            file_info = DagFileInfo(rel_path=rel_path, 
bundle_name=dag.bundle_name)
+            if file_info not in last_parsed:
+                # Zip-packaged dags are keyed by the archive path, not the 
inner file, so try the parent as well
+                file_info = DagFileInfo(rel_path=rel_path.parent, 
bundle_name=dag.bundle_name)
             if last_finish_time := last_parsed.get(file_info, None):
                 if dag.last_parsed_time + 
timedelta(seconds=self.stale_dag_threshold) < last_finish_time:
                     self.log.info(
diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py 
b/airflow-core/tests/unit/dag_processing/test_manager.py
index f1fb4b66c81..b136a4d9e8f 100644
--- a/airflow-core/tests/unit/dag_processing/test_manager.py
+++ b/airflow-core/tests/unit/dag_processing/test_manager.py
@@ -576,7 +576,6 @@ class TestDagFileProcessorManager:
         assert manager._file_queue == expected
 
         # Verify running it again produces same order
-        manager._files = []
         manager.prepare_file_queue(known_files=known_files)
         assert manager._file_queue == expected
 
@@ -1016,7 +1015,6 @@ class TestDagFileProcessorManager:
             run_count=1,
             last_num_of_db_queries=1,
         )
-        manager._files = [test_dag_path]
         manager._file_stats[test_dag_path] = stat
 
         active_dag_count = session.scalar(
@@ -1046,6 +1044,73 @@ class TestDagFileProcessorManager:
         # SerializedDagModel gives history about Dags
         assert serialized_dag_count == 1
 
+    @pytest.mark.usefixtures("testing_dag_bundle")
+    def test_scan_stale_dags_deactivates_zip_packaged_dags(self, session, 
test_zip_path):
+        """
+        Ensure that zip-packaged DAGs are marked inactive when the file is 
parsed but the
+        DagModel.last_parsed_time is not updated, testing fallback to the 
parent path when
+        comparing DAG.relative_fileloc to last_parsed entries.
+        """
+        manager = DagFileProcessorManager(
+            max_runs=1,
+            processor_timeout=10 * 60,
+        )
+        bundle = MagicMock()
+        bundle.name = "testing"
+        manager._dag_bundles = [bundle]
+
+        test_dag_path = DagFileInfo(
+            rel_path=Path("test_zip.zip"),
+            bundle_path=Path(test_zip_path).parent,
+            bundle_name="testing",
+        )
+        dagbag = DagBag(
+            test_dag_path.absolute_path,
+            bundle_path=test_dag_path.bundle_path,
+        )
+
+        # Add stale DAG to the DB
+        dag = dagbag.get_dag("test_zip_dag")
+        sync_dag_to_db(dag, session=session)
+
+        # Add DAG to the file_parsing_stats
+        stat = DagFileStat(
+            num_dags=1,
+            import_errors=0,
+            last_finish_time=timezone.utcnow() + timedelta(hours=1),
+            last_duration=1,
+            run_count=1,
+            last_num_of_db_queries=1,
+        )
+        manager._file_stats[test_dag_path] = stat
+
+        active_dag_count = session.scalar(
+            select(func.count(DagModel.dag_id)).where(
+                DagModel.dag_id == "test_zip_dag",
+                ~DagModel.is_stale,
+                DagModel.relative_fileloc == str(test_dag_path.rel_path / 
"test_zip.py"),
+            )
+        )
+        assert active_dag_count == 1
+
+        manager._scan_stale_dags()
+
+        active_dag_count = session.scalar(
+            select(func.count(DagModel.dag_id)).where(
+                DagModel.dag_id == "test_zip_dag",
+                ~DagModel.is_stale,
+                DagModel.relative_fileloc == str(test_dag_path.rel_path / 
"test_zip.py"),
+            )
+        )
+        assert active_dag_count == 0
+
+        serialized_dag_count = session.scalar(
+            
select(func.count(SerializedDagModel.dag_id)).where(SerializedDagModel.dag_id 
== dag.dag_id)
+        )
+        # Deactivating the DagModel should not delete the SerializedDagModel
+        # SerializedDagModel gives history about Dags
+        assert serialized_dag_count == 1
+
     @pytest.mark.usefixtures("testing_dag_bundle")
     def test_deactivate_stale_dags_marks_dags_in_inactive_bundles(self, 
session):
         """Dags whose bundle is no longer active should be marked stale even 
without a parse signal."""

Reply via email to