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


##########
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):

Review Comment:
   On main, `find_dag_file_paths` wrapped each file's `might_contain_dag` in a 
`try/except` and skipped the bad file. `PythonDagImporter.list_dag_definitions` 
calls `might_contain_dag` (and through it `read_bytes()`) with no guard, so a 
`PermissionError`, a file deleted mid-walk, or a `might_contain_dag_callable` 
that raises on one file ends the generator. The `except` at line 962 can't 
resume it, so the whole bundle is skipped on every refresh. I reproduced it 
with a bundle holding `good.py` plus a `chmod 000` file: main lists 
`['good.py']`, this branch raises `PermissionError`.
   
   `ZipImporter` already isolates unreadable members by yielding a 
`DagImportError`. Doing the same per file in 
`PythonDagImporter.list_dag_definitions` would queue the file, surface an 
import error for it, and keep the rest of the bundle. (#74029 keeps the 
per-file guard, for what it's worth.)



##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -950,10 +957,12 @@ def _refresh_dag_bundles(self, known_files: dict[str, 
set[DagFileInfo]]):
                 self._bundle_versions[bundle.name] = version_after_refresh
                 self._bundle_version_data[bundle.name] = 
version_data_after_refresh
 
-            found_files = {
-                DagFileInfo(rel_path=p, bundle_name=bundle.name, 
bundle_path=bundle.path)
-                for p in self._find_files_in_bundle(bundle)
-            }
+            try:
+                found_files = self._find_files_in_bundle(bundle)

Review Comment:
   `update_bundle_state(..., version=version_after_refresh)` and 
`self._bundle_versions[bundle.name] = version_after_refresh` (line 957) both 
run before this `try`, so by the time the listing raises the processor has 
already recorded the new version. On the next refresh `previously_seen` is true 
and the version hasn't moved, so the "version not changed" short-circuit above 
hits `continue` before we list again. For a git bundle that means the new files 
in that commit are never queued until another commit lands, and if it happens 
on the first refresh after startup nothing in the bundle gets parsed by this 
processor. A priority parse request doesn't get it out either, since 
`_force_refresh_bundles` only bypasses `should_skip_refresh`. On main the 
exception crashed the loop and the restart re-listed the bundle.
   
   Could the `except` pop `self._bundle_versions[bundle.name]` and 
`self._bundle_version_data[bundle.name]`, or could the version only be advanced 
once the listing succeeds? The comment at line 951 already says 
`_bundle_versions` is only advanced on success.



##########
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:
   The importer is discarded here and every definition it lists is kept, so an 
importer's `might_contain_dag` override only takes effect if its own 
`list_dag_definitions` calls it. Only `PythonDagImporter` and `ZipImporter` do. 
The `AbstractDagImporter.might_contain_dag` docstring tells importer authors 
that overriding it means "obvious non-DAG sources are dropped during discovery".
   
   I checked with a minimal `.yaml` importer that lists through 
`find_file_dag_definitions` and overrides only `might_contain_dag`: in safe 
mode a `not_a_dag.yaml` containing `foo: bar` is queued next to 
`pipeline.yaml`, so every non-Dag YAML file in a bundle gets its own parse 
process. #74029 gates each file through its importer's `might_contain_dag`, so 
the two PRs currently disagree on who owns this check. Should the gate live 
here (`for importer, item in ...`, skipping non-error items where `not 
importer.might_contain_dag(item, safe_mode)`), or should the base docstring say 
`list_dag_definitions` has to apply it itself? Gating here would sniff `.py` 
files twice, so the docstring route may be cleaner.



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

Review Comment:
   This is now the first place the manager builds a bundle's registry, and it 
runs inside the loop, after `before_run()` has called `gc.freeze()`. 
`list_dag_definitions` also materialises any configured importer specs, so a 
custom importer's module and its dependencies get imported here too. On Linux 
the parse child is a bare fork and `DagBag.__init__` picks up the same cached 
object from `get_importer_registry`, so the children inherit these objects 
outside the frozen set. Would it be worth warming the registry for each bundle 
(pending specs included) in `before_run` ahead of the freeze, so the parent's 
copy is the shared one?



##########
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:
   An item with no enclosing file is dropped here without a log. That covers a 
custom importer yielding a definition that isn't backed by a file (the 
`DagDefinition` ABC allows that) and a `DagImportError` whose 
`source_reference` isn't a bundle path, in which case the discovery error is 
lost as well. A `log.warning` in an `else` branch would tell the importer 
author why their Dags never show up.



##########
airflow-core/tests/unit/dag_processing/test_manager.py:
##########
@@ -3887,6 +3964,32 @@ def 
test_refresh_dag_bundles_update_bundle_state_failure_still_scans_files(self)
         # iteration will see a version mismatch and re-refresh rather than 
skip incorrectly
         assert "mock_bundle" not in manager._bundle_versions
 
+    def 
test_refresh_dag_bundles_discovery_failure_keeps_known_files_and_dags(self):

Review Comment:
   `_make_refresh_bundle()` defaults to `supports_versioning=False`, so this 
test never reaches the stuck-version path from the comment on line 961. 
Parametrizing over `supports_versioning` and asserting that a second 
`_refresh_dag_bundles` call lists again would pin that fix.



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

Review Comment:
   #74029 changes this same discovery path in `manager.py` and `utils/file.py`, 
with a different approach. Could the description link it, so the two can be 
reconciled before either merges? The comments on lines 991 and 996 cover where 
they currently differ (the per-file error guard, and who calls an importer's 
`might_contain_dag`).



##########
airflow-core/src/airflow/dag_processing/dagbag.py:
##########
@@ -276,9 +276,7 @@ def get_dag(self, dag_id, *, session: Session = 
NEW_SESSION):
             self.dags.pop(dag_id, None)
         if dag is None or is_expired:
             # Reprocess source file.
-            found_dags = self.process_file(
-                filepath=correct_maybe_zipped(orm_dag.fileloc), 
only_if_updated=False
-            )
+            found_dags = self.process_file(filepath=orm_dag.fileloc, 
only_if_updated=False)

Review Comment:
   Nothing pins this change. With `correct_maybe_zipped` restored here the 
related `get_dag` / refresh / zip tests all still pass, because 
`test_refresh_packaged_dag` only counts imports of `test_zip.py`, so 
re-importing the whole archive and re-importing one member look the same. 
Asserting that the other archive members are not re-imported would bind it.



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

Review Comment:
   Two user-visible changes come in through the registry here that the 
description doesn't mention:
   
   - Sourceless `.pyc` files are now queued for parsing. 
`PythonDagImporter.supported_extensions` is `[".py", ".pyc"]` and 
`find_file_dag_definitions` keeps a `.pyc` with no `.py` sibling, while the old 
scanner only took `.py` files and real zips.
   - A corrupt `.zip` is now queued with a `zip_read_error` import error from 
`ZipImporter.list_dag_definitions`. Before, `zipfile.is_zipfile` returned false 
and the file was skipped without a trace.
   
   Both look like reasonable behaviour, but users will see new import errors 
after upgrading, so they are worth a line in the description.



##########
airflow-core/src/airflow/dag_processing/manager.py:
##########
@@ -950,10 +957,12 @@ def _refresh_dag_bundles(self, known_files: dict[str, 
set[DagFileInfo]]):
                 self._bundle_versions[bundle.name] = version_after_refresh
                 self._bundle_version_data[bundle.name] = 
version_data_after_refresh
 
-            found_files = {
-                DagFileInfo(rel_path=p, bundle_name=bundle.name, 
bundle_path=bundle.path)
-                for p in self._find_files_in_bundle(bundle)
-            }
+            try:
+                found_files = self._find_files_in_bundle(bundle)
+            except Exception:
+                # Treating a failed listing as an empty bundle would 
deactivate all of its Dags.

Review Comment:
   This comment and the PR description both say a failed listing used to be 
treated as an empty bundle, but that isn't what main does. There 
`_find_files_in_bundle` ran without a `try`, and `_run_parsing_loop` calls 
`_refresh_dag_bundles` without one too, so an exception from the listing left 
the parsing loop. The only empty-bundle case on main was a missing bundle 
directory, where `list_py_file_paths` returns `[]`. Could the comment and the 
description describe what the `except` does now, rather than an old behaviour 
it replaces?



##########
airflow-core/tests/unit/utils/test_file.py:
##########
@@ -181,38 +180,6 @@ def test_get_modules_from_invalid_file(self):
 
         assert len(modules) == 0
 
-    def test_list_py_file_paths(self, test_zip_path):

Review Comment:
   This was the only test checking that discovery over a real folder honours 
`.airflowignore`. With `find_path_from_directory` in 
`airflow.sdk.importers.base` swapped for a plain `os.walk`, the core and 
task-sdk importer tests all still pass. A `_find_files_in_bundle` test with an 
`.airflowignore` in `tmp_path` (one ignored `.py`, one ignored `.zip`, one kept 
file) would cover it. Small thing in the same file: `TEST_DAG_FOLDER` (line 38) 
is now unused, and `TestListPyFilesPath` is still named after the removed 
function.



##########
airflow-core/src/airflow/utils/file.py:
##########
@@ -32,7 +32,6 @@
     might_contain_dag as might_contain_dag,
     might_contain_dag_via_default_heuristic as 
might_contain_dag_via_default_heuristic,
 )
-from airflow.configuration import conf
 
 log = logging.getLogger(__name__)

Review Comment:
   `import logging` and `log` have no users now that `find_dag_file_paths` is 
gone.



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