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


##########
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:
   Added one liner about the #74029 in the PR description - overall I think TP 
thought the discovery was forgotten before creating #74029, but I just divided 
it into 2 parts (wiring + this discovery). 



##########
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:
   Replaced the comment in the file and PR description with new behavior.



##########
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:
   Good catch, moved the listing before the version advance and now it only 
runs if listing succeeds + added test case for it.



##########
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:
   I modified `test_refresh_packaged_dag` now and it adds a second dag member 
and asserts that only the expired dag's member is re-imported. It fails if 
correct_maybe_zipped is restored.



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