bunnysayzz commented on code in PR #73979:
URL: https://github.com/apache/airflow/pull/73979#discussion_r4150073456


##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -1218,24 +1218,48 @@ def handle_removed_files(self, known_files: dict[str, 
set[DagFileInfo]]):
     def purge_removed_files_from_queue(self, present: set[DagFileInfo]):
         """Remove from queue any files no longer observed locally."""
         present_keys = {file.presence_key for file in present}
-        self._file_queue = OrderedDict((x, None) for x in self._file_queue if 
x.presence_key in present_keys)
+        self._file_queue = OrderedDict(
+            (x, None) for x in self._file_queue if self._file_is_present(x, 
present_keys)
+        )
         stats.gauge("dag_processing.file_path_queue_size", 
len(self._file_queue))
 
     def remove_orphaned_file_stats(self, present: set[DagFileInfo]):
         """Remove the stats for any dag files that don't exist anymore."""
         present_keys = {file.presence_key for file in present}
-        stats_to_remove = {file for file in self._file_stats if 
file.presence_key not in present_keys}
+        stats_to_remove = {
+            file for file in self._file_stats if not 
self._file_is_present(file, present_keys)
+        }
         for file in stats_to_remove:
             del self._file_stats[file]
 
+    @staticmethod
+    def _file_is_present(file: DagFileInfo, present_keys: set[tuple[str, 
Path]]) -> bool:
+        """
+        Check whether a tracked file is still observed in the bundle scan.
+
+        Dag files inside a zip archive are scanned as the archive itself
+        (``_find_files_in_bundle`` yields the ``.zip`` entry, not the inner
+        files), so an inner path counts as present while its containing
+        archive is still observed. Without this, a bundle refresh would
+        treat e.g. ``my_dags.zip/my_dag.py`` as removed and kill a
+        still-running callback processor for it.

Review Comment:
   done, moved the explanation above the def as a comment and kept the 
docstring to one line.



-- 
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]

Reply via email to