SameerMesiah97 commented on code in PR #73728:
URL: https://github.com/apache/airflow/pull/73728#discussion_r4108837785


##########
task-sdk/src/airflow/sdk/importers/base.py:
##########
@@ -483,6 +485,45 @@ def supported_extensions(self) -> list[str]:
         """Return all registered file extensions."""
         return sorted(set(self._extension_importers) | 
set(self._extension_specs))
 
+    def list_dag_definitions(
+        self,
+        bundle: BaseDagBundle,
+        *,
+        safe_mode: bool = True,
+    ) -> Iterator[tuple[AbstractDagImporter[Any], DagDefinition | 
DagImportError]]:
+        """
+        List DAG definitions in a bundle across all registered importers.
+
+        Each item is paired with the importer that listed it, which is the one 
to import it
+        with: a composite importer lists definitions no extension lookup would 
route back to
+        it, such as the members of an archive. A :class:`DagImportError` item 
is a
+        discovery-time failure rather than a source to import.
+        """
+        for spec in list(self._extension_specs.values()):
+            self._materialise_spec(spec)
+
+        for importer in self._ordered_importers:
+            for item in importer.list_dag_definitions(bundle, 
safe_mode=safe_mode):
+                # A file whose extension was claimed by a later registration 
belongs to that
+                # importer; yielding it here as well would import it twice.
+                if (
+                    isinstance(item, FilesystemDagDefinition)
+                    and self._extension_importers.get(item.suffix, importer) 
is not importer
+                ):
+                    continue
+                yield importer, item
+
+    def _materialise_spec(self, spec: _ImporterSpec) -> 
AbstractDagImporter[Any]:
+        """Instantiate a configured spec and take over every extension it was 
registered for."""
+        importer = self._instantiate_spec(spec)
+        for ext, registered in list(self._extension_specs.items()):
+            if registered is spec:
+                self._extension_importers[ext] = importer
+                del self._extension_specs[ext]
+        if importer not in self._ordered_importers:
+            self._ordered_importers.append(importer)
+        return importer

Review Comment:
   I think this can instantiate the same configured importer more than once 
when it supports multiple extensions. For example, if one spec is registered 
for both `.yaml` and `.yml`, `list(self._extension_specs.values())` contains 
that spec twice. The first call to `_materialise_spec` removes both entries, 
but the second spec is still in the snapshot. It then creates another importer 
and appends it to `_ordered_importers`, even though there are no extensions 
left for it to take over.
   
   That means `list_dag_definitions()` could call the importer twice and yield 
duplicate definitions. Could we materialise each spec only once and add a test 
for an importer registered with multiple extensions?
   



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