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


##########
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:
   This could trigger a ValueError is `path` is not inside `bundle.path`. I’d 
do something like
   
   ```python
   try:
       rel_path = Path(os.path.normpath(path)).relative_to(bundle.path)
   except ValueError:
       self.log.warning(
           "Ignoring %r listed in bundle %s: its file resolves outside the 
bundle",
           item,
           bundle.name,
       )
       continue
   ```



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