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


##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -970,56 +990,47 @@ 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 
None:
+                self.log.warning(
+                    "Ignoring %r listed in bundle %s: no file in the bundle 
holds it", item, bundle.name
+                )
+                continue
+            definition_locs[path.relative_to(bundle.path)].add(rel_fileloc)

Review Comment:
   Thanks, I moved it outside even before `find_enclosing_file` because it 
could cause an issue. 



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