dilnazanlid commented on code in PR #74020:
URL: https://github.com/apache/airflow/pull/74020#discussion_r4164452987


##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -970,56 +979,43 @@ def _refresh_dag_bundles(self, known_files: dict[str, 
set[DagFileInfo]]):
             self._resort_file_queue()
             self._add_new_files_to_queue(known_files=known_files)
 
-    def _find_files_in_bundle(self, bundle: BaseDagBundle) -> list[Path]:
-        """Get relative paths for dag files from bundle dir."""
-        # Build up a list of Python files that could contain DAGs
-        self.log.info("Searching for files in %s at %s", bundle.name, 
bundle.path)
-        rel_paths = [
-            Path(x).relative_to(bundle.path)
-            for x in list_py_file_paths(bundle.path, 
safe_mode=self.dag_discovery_safe_mode)
-        ]
+    def _find_files_in_bundle(self, bundle: BaseDagBundle) -> set[DagFileInfo]:
+        """
+        List the files to parse in a bundle through its importers.
+
+        A file holding several Dag definitions (a zip archive, for instance) 
is parsed as one.
+        """
+        self.log.info("Searching for Dag definitions in %s at %s", 
bundle.name, bundle.path)
+        registry = get_importer_registry(bundle.name)
+        definition_locs: defaultdict[Path, set[str]] = defaultdict(set)
+        for _, item in registry.list_dag_definitions(bundle, 
safe_mode=self.dag_discovery_safe_mode):
+            if isinstance(item, DagImportError):
+                # Importers report a source either absolutely or relative to 
the bundle.
+                rel_fileloc = os.path.relpath(bundle.path / 
item.source_reference, bundle.path)
+            else:
+                rel_fileloc = item.get_relative_loc(bundle.path)

Review Comment:
   I went with the docstring change - the registry can't apply the right check 
for composites anyway, since it pairs a zip member with ZipImporter while the 
meaningful check is the member importer's. Thanks!



##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -970,56 +979,43 @@ def _refresh_dag_bundles(self, known_files: dict[str, 
set[DagFileInfo]]):
             self._resort_file_queue()
             self._add_new_files_to_queue(known_files=known_files)
 
-    def _find_files_in_bundle(self, bundle: BaseDagBundle) -> list[Path]:
-        """Get relative paths for dag files from bundle dir."""
-        # Build up a list of Python files that could contain DAGs
-        self.log.info("Searching for files in %s at %s", bundle.name, 
bundle.path)
-        rel_paths = [
-            Path(x).relative_to(bundle.path)
-            for x in list_py_file_paths(bundle.path, 
safe_mode=self.dag_discovery_safe_mode)
-        ]
+    def _find_files_in_bundle(self, bundle: BaseDagBundle) -> set[DagFileInfo]:
+        """
+        List the files to parse in a bundle through its importers.
+
+        A file holding several Dag definitions (a zip archive, for instance) 
is parsed as one.
+        """
+        self.log.info("Searching for Dag definitions in %s at %s", 
bundle.name, bundle.path)
+        registry = get_importer_registry(bundle.name)
+        definition_locs: defaultdict[Path, set[str]] = defaultdict(set)
+        for _, item in registry.list_dag_definitions(bundle, 
safe_mode=self.dag_discovery_safe_mode):
+            if isinstance(item, DagImportError):
+                # Importers report a source either absolutely or relative to 
the bundle.
+                rel_fileloc = os.path.relpath(bundle.path / 
item.source_reference, bundle.path)
+            else:
+                rel_fileloc = item.get_relative_loc(bundle.path)
+            if (path := find_enclosing_file(bundle.path / rel_fileloc)) is not 
None:

Review Comment:
   Fixed, added warning log. (although this should become obsolete with the 
#73457 where the files are fully replaced by dag definition objects)



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