This is an automated email from the ASF dual-hosted git repository.

uranusjr pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new adaa16f88af Parse Dag files through the Task SDK importers (#73853)
adaa16f88af is described below

commit adaa16f88af55076e19dc13f131dde34a35ab5c6
Author: Dilnaz Amanzholova <[email protected]>
AuthorDate: Thu Oct 1 11:42:47 2026 +0200

    Parse Dag files through the Task SDK importers (#73853)
---
 .../src/airflow/config_templates/config.yml        |  37 +++
 .../src/airflow/dag_processing/collection.py       |  20 +-
 airflow-core/src/airflow/dag_processing/dagbag.py  | 232 ++++++-------
 .../airflow/dag_processing/importers/__init__.py   |  39 ---
 .../src/airflow/dag_processing/importers/base.py   | 266 ---------------
 .../dag_processing/importers/python_importer.py    | 360 ---------------------
 airflow-core/src/airflow/dag_processing/manager.py |   2 +
 .../src/airflow/dag_processing/processor.py        |  15 +-
 airflow-core/src/airflow/models/dagcode.py         |  45 ++-
 airflow-core/src/airflow/models/serialized_dag.py  |  19 +-
 airflow-core/src/airflow/utils/file.py             |  10 +
 .../unit/dag_processing/importers/__init__.py      |  16 -
 .../unit/dag_processing/importers/test_registry.py |  99 ------
 .../tests/unit/dag_processing/test_collection.py   |  52 +++
 .../tests/unit/dag_processing/test_dagbag.py       | 263 +++++++++------
 .../tests/unit/dag_processing/test_manager.py      |  25 ++
 .../tests/unit/dag_processing/test_processor.py    |  31 +-
 airflow-core/tests/unit/models/test_dagcode.py     |  39 +++
 devel-common/src/tests_common/test_utils/config.py |  11 +-
 generated/known_sdk_imports_in_core.txt            |   6 +-
 .../airflow/sdk/execution_time/schema/schema.json  |  33 ++
 .../sdk/execution_time/schema/versions/__init__.py |   8 +-
 .../execution_time/schema/versions/v2026_10_30.py  |  12 +
 task-sdk/src/airflow/sdk/importers/base.py         |   2 +
 ts-sdk/src/generated/supervisor.ts                 |  18 ++
 25 files changed, 648 insertions(+), 1012 deletions(-)

diff --git a/airflow-core/src/airflow/config_templates/config.yml 
b/airflow-core/src/airflow/config_templates/config.yml
index 4d3677a1d4a..cc8f4c54c05 100644
--- a/airflow-core/src/airflow/config_templates/config.yml
+++ b/airflow-core/src/airflow/config_templates/config.yml
@@ -3223,6 +3223,10 @@ dag_processor:
 
         The default is the dags folder dag bundle.
 
+        A backend config may also supply ``importers`` (experimental): Dag 
importers for that
+        bundle only, in the format of ``[dag_processor] 
dag_importer_configs``, taking precedence
+        over it.
+
         Note: As shown below, you can split your json config over multiple 
lines by indenting.
         See configparser documentation for an example:
         
https://docs.python.org/3/library/configparser.html#supported-ini-file-structure.
@@ -3248,6 +3252,39 @@ dag_processor:
                 "kwargs": {{}}
               }}
             ]
+
+    dag_importer_configs:
+      description: |
+        .. note:: |experimental|
+
+        List of Dag importer configs. Each entry must supply ``classpath`` and 
may supply
+        ``extensions`` and ``kwargs``.
+
+        Importers turn the contents of a bundle into Dag definitions. Airflow 
registers
+        importers for ``.py`` and ``.zip`` out of the box; use this option to 
register
+        importers for other formats, or to replace a built-in one.
+
+        An importer registered here applies to every bundle. To register an 
importer for a
+        single bundle only, add an ``importers`` key to that bundle's entry in
+        ``dag_bundle_config_list``, which takes the same format and takes 
precedence over
+        this option.
+
+        If ``extensions`` is omitted, the extensions declared by the importer 
class are used.
+
+        An entry may also be given as a bare classpath string when it needs no 
other
+        settings, so ``["my_company.importers.YamlDagImporter"]`` is 
equivalent to
+        ``[{"classpath": "my_company.importers.YamlDagImporter"}]``.
+      version_added: 3.4.0
+      type: string
+      example: >
+        [
+              {
+                "classpath": "my_company.importers.YamlDagImporter",
+                "extensions": [".yaml", ".yml"],
+                "kwargs": {}
+              }
+            ]
+      default: ~
     refresh_interval:
       description: |
         How often (in seconds) to refresh, or look for new files, in a DAG 
bundle.
diff --git a/airflow-core/src/airflow/dag_processing/collection.py 
b/airflow-core/src/airflow/dag_processing/collection.py
index 46099badce7..ec8889db318 100644
--- a/airflow-core/src/airflow/dag_processing/collection.py
+++ b/airflow-core/src/airflow/dag_processing/collection.py
@@ -31,7 +31,7 @@ import traceback
 from typing import TYPE_CHECKING, Any, NamedTuple, TypeVar
 
 import structlog
-from sqlalchemy import delete, false, func, insert, select, tuple_, update
+from sqlalchemy import delete, false, func, insert, or_, select, tuple_, update
 from sqlalchemy.exc import OperationalError
 from sqlalchemy.orm import joinedload, load_only
 
@@ -78,6 +78,7 @@ if TYPE_CHECKING:
     from sqlalchemy.sql import Select
 
     from airflow.models.serialized_dag import DagWriteMetadata
+    from airflow.sdk.importers import DagSourceCode  # noqa: SDK001
     from airflow.typing_compat import Self, Unpack
 
     AssetT = TypeVar("AssetT", SerializedAsset, SerializedAssetAlias)
@@ -259,6 +260,7 @@ def _serialize_dag_capturing_errors(
     bundle_version: str | None,
     version_data: dict | None = None,
     _prefetched: DagWriteMetadata | None = None,
+    dag_source_code: DagSourceCode | None = None,
 ):
     """
     Try to serialize the dag to the DB, but make a note of any errors.
@@ -280,12 +282,15 @@ def _serialize_dag_capturing_errors(
             bundle_version=bundle_version,
             version_data=version_data,
             min_update_interval=MIN_SERIALIZED_DAG_UPDATE_INTERVAL,
+            dag_source_code=dag_source_code,
             session=session,
             _prefetched=_prefetched,
         )
         if not dag_was_updated:
             # Check and update DagCode
-            DagCode.update_source_code(dag.dag_id, dag.fileloc, 
session=session)
+            DagCode.update_source_code(
+                dag.dag_id, dag.fileloc, dag_source_code=dag_source_code, 
session=session
+            )
         if "FabAuthManager" in conf.get("core", "auth_manager"):
             _sync_dag_perms(dag, session=session)
 
@@ -370,7 +375,11 @@ def _update_dag_warnings(
         session.scalars(
             select(DagWarning).where(
                 DagWarning.dag_id.in_(dag_ids),
-                DagWarning.warning_type.in_(warning_types),
+                or_(
+                    DagWarning.warning_type.in_(warning_types),
+                    # Importer-namespaced types are only ever reported while 
parsing.
+                    DagWarning.warning_type.not_in([t.value for t in 
DagWarningType]),
+                ),
             )
         )
     )
@@ -588,6 +597,7 @@ def update_dag_parsing_results_in_db(
         DagWarningType.RUNTIME_VARYING_VALUE,
     ),
     files_parsed: set[tuple[str, str]] | None = None,
+    dag_source_codes: dict[str, DagSourceCode] | None = None,
 ):
     """
     Update everything to do with DAG parsing in the DB.
@@ -609,6 +619,8 @@ def update_dag_parsing_results_in_db(
     :param files_parsed: Set of (bundle_name, relative_fileloc) tuples for all 
files that were parsed.
         If None, will be inferred from dags and import_errors. Passing this 
explicitly ensures that
         import errors are cleared for files that were parsed but no longer 
contain DAGs.
+    :param dag_source_codes: Source code read by the Dag importers, keyed by 
Dag fileloc. Dags
+        without an entry have their source read from ``fileloc``.
     """
     accepted = _reject_other_teams_plugin_classes(bundle_name, dags, 
import_errors, session=session)
     if len(accepted) != len(dags):
@@ -617,6 +629,7 @@ def update_dag_parsing_results_in_db(
         warnings = {warning for warning in warnings if warning.dag_id not in 
rejected_ids}
     dags = accepted
 
+    dag_source_codes = dag_source_codes or {}
     # Retry 'DAG.bulk_write_to_db' & 'SerializedDagModel.bulk_sync_to_db' in 
case
     # of any Operational Errors
     # In case of failures, provide_session handles rollback
@@ -657,6 +670,7 @@ def update_dag_parsing_results_in_db(
                             version_data=version_data,
                             session=session,
                             _prefetched=prefetched_metadata.get(dag.dag_id),
+                            dag_source_code=dag_source_codes.get(dag.fileloc),
                         )
                     )
             except OperationalError:
diff --git a/airflow-core/src/airflow/dag_processing/dagbag.py 
b/airflow-core/src/airflow/dag_processing/dagbag.py
index 584840f25c3..295eb52536b 100644
--- a/airflow-core/src/airflow/dag_processing/dagbag.py
+++ b/airflow-core/src/airflow/dag_processing/dagbag.py
@@ -22,8 +22,8 @@ import os
 import sys
 import textwrap
 import warnings
-from collections.abc import Generator
-from datetime import datetime, timedelta
+from collections.abc import Iterator
+from datetime import timedelta
 from pathlib import Path
 from typing import TYPE_CHECKING, Any, NamedTuple
 
@@ -32,7 +32,7 @@ from tabulate import tabulate
 from airflow import settings
 from airflow._shared.timezones import timezone
 from airflow.configuration import conf
-from airflow.dag_processing.importers import get_importer_registry
+from airflow.dag_processing.bundles.local import LocalDagBundle
 from airflow.exceptions import (
     AirflowClusterPolicyError,
     AirflowClusterPolicySkipDag,
@@ -44,9 +44,10 @@ from airflow.exceptions import (
 from airflow.executors.executor_loader import ExecutorLoader
 from airflow.listeners.listener import get_listener_manager
 from airflow.models.pool import Pool
+from airflow.sdk.importers import DagImportError, get_importer_registry
 from airflow.serialization.definitions.notset import NOTSET, ArgNotSet, 
is_arg_set
 from airflow.serialization.serialized_objects import LazyDeserializedDAG
-from airflow.utils.file import correct_maybe_zipped
+from airflow.utils.file import correct_maybe_zipped, find_enclosing_file
 from airflow.utils.log.logging_mixin import LoggingMixin
 from airflow.utils.session import NEW_SESSION, provide_session
 
@@ -55,25 +56,7 @@ if TYPE_CHECKING:
 
     from airflow import DAG
     from airflow.models.dagwarning import DagWarning
-
-
[email protected]
-def _capture_with_reraise() -> Generator[list[warnings.WarningMessage], None, 
None]:
-    """Capture warnings in context and re-raise it on exit from the context 
manager."""
-    captured_warnings = []
-    try:
-        with warnings.catch_warnings(record=True) as captured_warnings:
-            yield captured_warnings
-    finally:
-        if captured_warnings:
-            for cw in captured_warnings:
-                warnings.warn_explicit(
-                    message=cw.message,
-                    category=cw.category,
-                    filename=cw.filename,
-                    lineno=cw.lineno,
-                    source=cw.source,
-                )
+    from airflow.sdk.importers import AbstractDagImporter, DagDefinition, 
DagImportWarning, DagSourceCode
 
 
 class FileLoadStat(NamedTuple):
@@ -224,16 +207,24 @@ class DagBag(LoggingMixin):
         dag_folder = dag_folder or settings.DAGS_FOLDER
         self.dag_folder = dag_folder
         self.dags: dict[str, DAG] = {}
-        # the file's last modified timestamp when we last read it
-        self.file_last_changed: dict[str, datetime] = {}
+        # The freshness token of each definition when we last imported it, 
keyed by its fileloc
+        self.file_last_changed: dict[str, str] = {}
         # Store import errors with relative file paths as keys (relative to 
bundle_path)
         self.import_errors: dict[str, str] = {}
         self.captured_warnings: dict[str, tuple[str, ...]] = {}
-        self.has_logged = False
+        self._import_warnings: dict[str, list[DagImportWarning]] = {}
+        # The source code of each definition that produced Dags, keyed by its 
fileloc
+        self.dag_source_codes: dict[str, DagSourceCode] = {}
         # Only used by SchedulerJob to compare the dag_hash to identify change 
in DAGs
         self.dags_hash: dict[str, str] = {}
 
         self.known_pools = known_pools
+        self._importer_registry = get_importer_registry(bundle_name)
+        # Importers only read the bundle's name and path, so a bag built 
without a bundle
+        # hands them its Dag folder instead.
+        self._bundle = LocalDagBundle(
+            name=bundle_name or "dags-folder", path=os.fspath(bundle_path or 
dag_folder)
+        )
 
         self.dagbag_import_error_tracebacks = conf.getboolean("core", 
"dagbag_import_error_tracebacks")
         self.dagbag_import_error_traceback_depth = conf.getint("core", 
"dagbag_import_error_traceback_depth")
@@ -299,74 +290,74 @@ class DagBag(LoggingMixin):
         """Process a DAG file and return found DAGs."""
         if filepath is None or not os.path.isfile(filepath):
             return []
+        return [
+            dag
+            for importer, definition in self._find_definitions(Path(filepath), 
safe_mode=safe_mode)
+            for dag in self._process_definition(importer, definition, 
only_if_updated=only_if_updated)
+        ]
+
+    def _find_definitions(
+        self, path: Path, *, safe_mode: bool
+    ) -> Iterator[tuple[AbstractDagImporter, DagDefinition]]:
+        """
+        Discover the Dag definitions at ``path`` with the importer that should 
import each.
 
-        try:
-            file_last_changed_on_disk = 
datetime.fromtimestamp(os.path.getmtime(filepath))
-            if (
-                only_if_updated
-                and filepath in self.file_last_changed
-                and file_last_changed_on_disk == 
self.file_last_changed[filepath]
-            ):
-                return []
-        except Exception as e:
-            self.log.exception(e)
-            return []
-
-        self.captured_warnings.pop(filepath, None)
-
-        registry = get_importer_registry()
-        importer = registry.get_importer(filepath)
-
-        if importer is None:
-            self.log.debug("No importer found for file: %s", filepath)
+        ``path`` may be a directory, a file, or a definition nested in a file 
(a zip member);
+        discovery is scoped to the directory or the enclosing file, and 
narrowed to the
+        definition itself in the nested case. Discovery errors are recorded as 
import errors.
+        """
+        scope = path if path.is_dir() else find_enclosing_file(path)
+        if scope is None:
+            return
+        scoped_bundle = LocalDagBundle(name=self._bundle.name, 
path=os.fspath(scope))
+        for importer, item in self._importer_registry.list_dag_definitions(
+            scoped_bundle, safe_mode=safe_mode
+        ):
+            if isinstance(item, DagImportError):
+                self._record_import_error(item, root=scope)
+            elif scope == path or Path(repr(item)) == path:
+                yield importer, item
+
+    def _process_definition(
+        self, importer: AbstractDagImporter, definition: DagDefinition, *, 
only_if_updated: bool
+    ) -> list[DAG]:
+        """Import a Dag definition and bag the Dags it defines."""
+        fileloc = repr(definition)
+        freshness_token = definition.freshness_token
+        if only_if_updated and self.file_last_changed.get(fileloc) == 
freshness_token:
             return []
 
-        result = importer.import_file(
-            file_path=filepath,
-            bundle_path=self.bundle_path,
-            bundle_name=self.bundle_name,
-            safe_mode=safe_mode,
-        )
+        self.captured_warnings.pop(fileloc, None)
+        self._import_warnings.pop(fileloc, None)
+        self.dag_source_codes.pop(fileloc, None)
+        result = importer.import_definition(definition, self._bundle)
 
-        if result.skipped_files:
-            for skipped in result.skipped_files:
-                if not self.has_logged:
-                    self.has_logged = True
-                    self.log.info("File %s assumed to contain no DAGs. 
Skipping.", skipped)
-
-        if result.errors:
-            for error in result.errors:
-                # Use the relative file path from error (importer provides 
relative paths)
-                # Fall back to converting filepath to relative if 
error.file_path is not set
-                error_path = error.file_path if error.file_path else 
self._get_relative_fileloc(filepath)
-                error_msg = error.stacktrace if error.stacktrace else 
error.message
-                self.import_errors[error_path] = error_msg
-                self.log.error("Error loading DAG from %s: %s", error_path, 
error.message)
+        for error in result.errors:
+            self._record_import_error(error, root=self._bundle.path)
 
         if result.warnings:
-            formatted_warnings = [
-                f"{w.file_path}:{w.line_number}: {w.warning_type}: 
{w.message}" for w in result.warnings
-            ]
-            self.captured_warnings[filepath] = tuple(formatted_warnings)
+            self._import_warnings[fileloc] = result.warnings
+            self.captured_warnings[fileloc] = tuple(
+                f"{w.source_reference}:{w.line_number}: {w.warning_type}: 
{w.message}"
+                for w in result.warnings
+            )
             # Re-emit warnings so they can be handled by Python's warning 
system
             for w in result.warnings:
                 warnings.warn_explicit(
                     message=w.message,
                     category=UserWarning,
-                    filename=w.file_path,
+                    filename=w.source_reference,
                     lineno=w.line_number or 0,
                 )
 
         bagged_dags = []
         for dag in result.dags:
             try:
-                if dag.fileloc is None:
-                    dag.fileloc = filepath
-
-                # Add the bundle_name to the Dag
+                # Importers are not required to set these, and ``DAG.fileloc`` 
otherwise
+                # defaults to the file that constructed the Dag (the 
importer's own module).
+                dag.fileloc = fileloc
+                dag.relative_fileloc = self._get_relative_fileloc(fileloc)
                 dag.bundle_name = self.bundle_name
-
-                # Validate before adding to bag (matches original 
_process_modules behavior)
                 dag.validate()
                 _validate_executor_fields(dag, self.bundle_name)
                 _assign_default_team_pools(dag, self.bundle_name)
@@ -375,51 +366,69 @@ class DagBag(LoggingMixin):
             except AirflowClusterPolicySkipDag:
                 self.log.debug("DAG %s skipped by cluster policy", dag.dag_id)
             except Exception as e:
-                self.log.exception("Error bagging DAG from %s", filepath)
-                relative_path = self._get_relative_fileloc(filepath)
-                self.import_errors[relative_path] = f"{type(e).__name__}: {e}"
+                self.log.exception("Error bagging DAG from %s", fileloc)
+                self.import_errors[self._get_relative_fileloc(fileloc)] = 
f"{type(e).__name__}: {e}"
 
-        self.file_last_changed[filepath] = file_last_changed_on_disk
+        self.file_last_changed[fileloc] = freshness_token
+        if bagged_dags:
+            try:
+                self.dag_source_codes[fileloc] = 
importer.get_source_code(definition)
+            except Exception:
+                self.log.exception("Failed to read source code of %s", fileloc)
         return bagged_dags
 
+    def _record_import_error(self, error: DagImportError, *, root: Path) -> 
None:
+        # Importers report a source either absolutely or relative to the 
bundle they were given.
+        fileloc = self._get_relative_fileloc(os.fspath(root / 
error.source_reference))
+        self.import_errors[fileloc] = error.stacktrace or error.message
+        self.log.error("Error loading DAG from %s: %s", fileloc, error.message)
+
+    @property
+    def parsed_definitions(self) -> list[str]:
+        """Locations of the Dag definitions imported into this bag, relative 
to the bundle."""
+        return [self._get_relative_fileloc(fileloc) for fileloc in 
self.file_last_changed]
+
     @property
     def dag_warnings(self) -> set[DagWarning]:
         """Get the set of DagWarnings for the bagged dags."""
         from airflow.models.dagwarning import DagWarning, DagWarningType
 
+        dag_warnings: set[DagWarning] = set()
+        for dag in self.dags.values():
+            for import_warning in self._import_warnings.get(dag.fileloc, ()):
+                # Only importer-namespaced types (``yaml:deprecated_field``) 
are Dag warnings;
+                # Python warning categories stay in ``captured_warnings``.
+                with contextlib.suppress(ValueError):
+                    dag_warnings.add(
+                        DagWarning(dag.dag_id, import_warning.warning_type, 
import_warning.message)
+                    )
+
         # None means this feature is not enabled. Empty set means we don't 
know about any pools at all!
         if self.known_pools is None:
-            return set()
+            return dag_warnings
 
-        def get_pools(dag) -> dict[str, set[str]]:
-            return {dag.dag_id: {task.pool for task in dag.tasks}}
-
-        pool_dict: dict[str, set[str]] = {}
         for dag in self.dags.values():
-            pool_dict.update(get_pools(dag))
-
-        warnings: set[DagWarning] = set()
-        for dag_id, dag_pools in pool_dict.items():
-            nonexistent_pools = dag_pools - self.known_pools
+            nonexistent_pools = {task.pool for task in dag.tasks} - 
self.known_pools
             if nonexistent_pools:
-                warnings.add(
+                dag_warnings.add(
                     DagWarning(
-                        dag_id,
+                        dag.dag_id,
                         DagWarningType.NONEXISTENT_POOL,
-                        f"Dag '{dag_id}' references non-existent pools: 
{sorted(nonexistent_pools)!r}",
+                        f"Dag '{dag.dag_id}' references non-existent pools: 
{sorted(nonexistent_pools)!r}",
                     )
                 )
-        return warnings
+        return dag_warnings
 
     def _get_relative_fileloc(self, filepath: str) -> str:
         """
         Get the relative file location for a given filepath.
 
         :param filepath: Absolute path to the file
-        :return: Relative path from bundle_path, or original filepath if no 
bundle_path
+        :return: Relative path from bundle_path, or original filepath if not 
under bundle_path
         """
         if self.bundle_path:
-            return str(Path(filepath).relative_to(self.bundle_path))
+            with contextlib.suppress(ValueError):
+                return str(Path(filepath).relative_to(self.bundle_path))
         return filepath
 
     def bag_dag(self, dag: DAG):
@@ -492,23 +501,18 @@ class DagBag(LoggingMixin):
         # Used to store stats around DagBag processing
         stats = []
 
-        # Ensure dag_folder is a str -- it may have been a pathlib.Path
-        dag_folder = correct_maybe_zipped(str(dag_folder))
-
-        registry = get_importer_registry()
-        files_to_parse = registry.list_dag_files(dag_folder, 
safe_mode=safe_mode)
-
-        for filepath in files_to_parse:
+        for importer, definition in self._find_definitions(Path(dag_folder), 
safe_mode=safe_mode):
+            fileloc = repr(definition)
             try:
                 file_parse_start_dttm = timezone.utcnow()
-                found_dags = self.process_file(filepath, 
only_if_updated=only_if_updated, safe_mode=safe_mode)
+                found_dags = self._process_definition(importer, definition, 
only_if_updated=only_if_updated)
 
                 file_parse_end_dttm = timezone.utcnow()
                 try:
-                    relative_file = 
Path(filepath).relative_to(Path(self.dag_folder)).as_posix()
+                    relative_file = 
Path(fileloc).relative_to(Path(self.dag_folder)).as_posix()
                 except ValueError:
-                    # filepath is not under dag_folder (e.g., example DAGs 
from a different location)
-                    relative_file = Path(filepath).as_posix()
+                    # fileloc is not under dag_folder (e.g., example DAGs from 
a different location)
+                    relative_file = Path(fileloc).as_posix()
                 stats.append(
                     FileLoadStat(
                         file=relative_file,
@@ -516,7 +520,7 @@ class DagBag(LoggingMixin):
                         dag_num=len(found_dags),
                         task_num=sum(len(dag.tasks) for dag in found_dags),
                         dags=str([dag.dag_id for dag in found_dags]),
-                        warning_num=len(self.captured_warnings.get(filepath, 
[])),
+                        warning_num=len(self.captured_warnings.get(fileloc, 
[])),
                         bundle_path=self.bundle_path,
                         bundle_name=self.bundle_name,
                     )
@@ -587,13 +591,14 @@ def sync_bag_to_db(
     import_errors = {(bundle_name, rel_path): error for rel_path, error in 
dagbag.import_errors.items()}
 
     # Build the set of all files that were parsed and include files with 
import errors
-    # in case they are not in file_last_changed
+    # in case they are not in parsed_definitions
     files_parsed = set(import_errors)
     if dagbag.bundle_path:
-        files_parsed.update(
-            (bundle_name, dagbag._get_relative_fileloc(abs_filepath))
-            for abs_filepath in dagbag.file_last_changed
-        )
+        for rel_path in dagbag.parsed_definitions:
+            files_parsed.add((bundle_name, rel_path))
+            # A definition nested in an archive also clears the archive's own 
discovery errors.
+            if enclosing_file := find_enclosing_file(Path(dagbag.bundle_path, 
rel_path)):
+                files_parsed.add((bundle_name, 
dagbag._get_relative_fileloc(os.fspath(enclosing_file))))
 
     update_dag_parsing_results_in_db(
         bundle_name,
@@ -605,4 +610,5 @@ def sync_bag_to_db(
         session=session,
         version_data=version_data,
         files_parsed=files_parsed,
+        dag_source_codes=dagbag.dag_source_codes,
     )
diff --git a/airflow-core/src/airflow/dag_processing/importers/__init__.py 
b/airflow-core/src/airflow/dag_processing/importers/__init__.py
deleted file mode 100644
index 15fe3734bc5..00000000000
--- a/airflow-core/src/airflow/dag_processing/importers/__init__.py
+++ /dev/null
@@ -1,39 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements.  See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership.  The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License.  You may obtain a copy of the License at
-#
-#   http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied.  See the License for the
-# specific language governing permissions and limitations
-# under the License.
-"""DAG Importers - pluggable mechanism for importing DAGs from different file 
formats."""
-
-from __future__ import annotations
-
-from airflow.dag_processing.importers.base import (
-    AbstractDagImporter,
-    DagImporterRegistry,
-    DagImportError,
-    DagImportResult,
-    DagImportWarning,
-    get_importer_registry,
-)
-from airflow.dag_processing.importers.python_importer import PythonDagImporter
-
-__all__ = [
-    "AbstractDagImporter",
-    "DagImportError",
-    "DagImporterRegistry",
-    "DagImportResult",
-    "DagImportWarning",
-    "PythonDagImporter",
-    "get_importer_registry",
-]
diff --git a/airflow-core/src/airflow/dag_processing/importers/base.py 
b/airflow-core/src/airflow/dag_processing/importers/base.py
deleted file mode 100644
index 036dfab98d7..00000000000
--- a/airflow-core/src/airflow/dag_processing/importers/base.py
+++ /dev/null
@@ -1,266 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements.  See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership.  The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License.  You may obtain a copy of the License at
-#
-#   http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied.  See the License for the
-# specific language governing permissions and limitations
-# under the License.
-"""Abstract base class for DAG importers."""
-
-from __future__ import annotations
-
-import logging
-import os
-import threading
-from abc import ABC, abstractmethod
-from collections.abc import Iterator
-from dataclasses import dataclass, field
-from pathlib import Path
-from typing import TYPE_CHECKING
-
-from airflow._shared.module_loading.file_discovery import 
find_path_from_directory
-from airflow.configuration import conf
-from airflow.utils.file import might_contain_dag
-
-if TYPE_CHECKING:
-    from airflow.sdk import DAG
-
-log = logging.getLogger(__name__)
-
-
-@dataclass
-class DagImportError:
-    """Structured error information for DAG import failures."""
-
-    file_path: str
-    message: str
-    error_type: str = "import"
-    line_number: int | None = None
-    column_number: int | None = None
-    context: str | None = None
-    suggestion: str | None = None
-    stacktrace: str | None = None
-
-    def format_message(self) -> str:
-        """Format the error as a human-readable string."""
-        parts = [f"Error in {self.file_path}"]
-        if self.line_number is not None:
-            loc = f"line {self.line_number}"
-            if self.column_number is not None:
-                loc += f", column {self.column_number}"
-            parts.append(f"Location: {loc}")
-        parts.append(f"Error ({self.error_type}): {self.message}")
-        if self.context:
-            parts.append(f"Context:\n{self.context}")
-        if self.suggestion:
-            parts.append(f"Suggestion: {self.suggestion}")
-        return "\n".join(parts)
-
-
-@dataclass
-class DagImportWarning:
-    """Warning information for non-fatal issues during DAG import."""
-
-    file_path: str
-    message: str
-    warning_type: str = "general"
-    line_number: int | None = None
-
-
-@dataclass
-class DagImportResult:
-    """Result of importing DAGs from a file."""
-
-    file_path: str
-    dags: list[DAG] = field(default_factory=list)
-    errors: list[DagImportError] = field(default_factory=list)
-    skipped_files: list[str] = field(default_factory=list)
-    warnings: list[DagImportWarning] = field(default_factory=list)
-
-    @property
-    def success(self) -> bool:
-        """Return True if no fatal errors occurred."""
-        return len(self.errors) == 0
-
-
-class AbstractDagImporter(ABC):
-    """Abstract base class for DAG importers."""
-
-    @classmethod
-    @abstractmethod
-    def supported_extensions(cls) -> list[str]:
-        """Return file extensions this importer handles (e.g., ['.py', 
'.zip'])."""
-
-    @abstractmethod
-    def import_file(
-        self,
-        file_path: str | Path,
-        *,
-        bundle_path: Path | None = None,
-        bundle_name: str | None = None,
-        safe_mode: bool = True,
-    ) -> DagImportResult:
-        """Import DAGs from a file."""
-
-    def can_handle(self, file_path: str | Path) -> bool:
-        """Check if this importer can handle the given file."""
-        path = Path(file_path) if isinstance(file_path, str) else file_path
-        return path.suffix.lower() in self.supported_extensions()
-
-    def get_relative_path(self, file_path: str | Path, bundle_path: Path | 
None) -> str:
-        """Get the relative file path from the bundle root."""
-        if bundle_path is None:
-            return str(file_path)
-        try:
-            return str(Path(file_path).relative_to(bundle_path))
-        except ValueError:
-            return str(file_path)
-
-    def list_dag_files(
-        self,
-        directory: str | os.PathLike[str],
-        safe_mode: bool = True,
-    ) -> Iterator[str]:
-        """
-        List DAG files in a directory that this importer can handle.
-
-        Override this method to customize file discovery for your importer.
-        The default implementation finds files matching supported_extensions()
-        and respects .airflowignore files.
-
-        :param directory: Directory to search for DAG files
-        :param safe_mode: Whether to use heuristics to filter non-DAG files
-        :return: Iterator of file paths
-        """
-        ignore_file_syntax = conf.get_mandatory_value("core", 
"DAG_IGNORE_FILE_SYNTAX", fallback="glob")
-        supported_exts = [ext.lower() for ext in self.supported_extensions()]
-
-        for file_path in find_path_from_directory(directory, ".airflowignore", 
ignore_file_syntax):
-            path = Path(file_path)
-
-            if not path.is_file():
-                continue
-
-            # Check if this importer handles this file extension
-            if path.suffix.lower() not in supported_exts:
-                continue
-
-            # Apply safe_mode heuristic if enabled
-            if safe_mode and not might_contain_dag(file_path, safe_mode, 
conf=conf):
-                continue
-
-            yield file_path
-
-
-class DagImporterRegistry:
-    """
-    Registry for DAG importers. Singleton that manages importers by file 
extension.
-
-    Each file extension can only be handled by one importer at a time. If 
multiple
-    importers claim the same extension, the last registered one wins and a 
warning
-    is logged. The built-in PythonDagImporter handles .py and .zip extensions.
-    """
-
-    _instance: DagImporterRegistry | None = None
-    _importers: dict[str, AbstractDagImporter]
-    _lock = threading.Lock()
-
-    def __new__(cls) -> DagImporterRegistry:
-        with cls._lock:
-            if cls._instance is None:
-                cls._instance = super().__new__(cls)
-                cls._instance._importers = {}
-                cls._instance._register_default_importers()
-        return cls._instance
-
-    def _register_default_importers(self) -> None:
-        from airflow.dag_processing.importers.python_importer import 
PythonDagImporter
-
-        self.register(PythonDagImporter())
-
-    def register(self, importer: AbstractDagImporter) -> None:
-        """
-        Register an importer for its supported extensions.
-
-        Each extension can only have one importer. If an extension is already 
registered,
-        the new importer will override it and a warning will be logged.
-        """
-        for ext in importer.supported_extensions():
-            ext_lower = ext.lower()
-            if ext_lower in self._importers:
-                existing = self._importers[ext_lower]
-                log.warning(
-                    "Extension '%s' already registered by %s, overriding with 
%s",
-                    ext,
-                    type(existing).__name__,
-                    type(importer).__name__,
-                )
-            self._importers[ext_lower] = importer
-
-    def get_importer(self, file_path: str | Path) -> AbstractDagImporter | 
None:
-        """Get the appropriate importer for a file, or None if unsupported."""
-        path = Path(file_path) if isinstance(file_path, str) else file_path
-        return self._importers.get(path.suffix.lower())
-
-    def can_handle(self, file_path: str | Path) -> bool:
-        """Check if any registered importer can handle this file."""
-        return self.get_importer(file_path) is not None
-
-    def supported_extensions(self) -> list[str]:
-        """Return all registered file extensions."""
-        return list(self._importers.keys())
-
-    def list_dag_files(
-        self,
-        directory: str | os.PathLike[str],
-        safe_mode: bool = True,
-    ) -> list[str]:
-        """
-        List all DAG files in a directory using all registered importers.
-
-        If directory is actually a file, returns that file if any importer can 
handle it.
-
-        :param directory: Directory (or file) to search for DAG files
-        :param safe_mode: Whether to use heuristics to filter non-DAG files
-        :return: List of file paths (deduplicated)
-        """
-        path = Path(directory)
-
-        # If it's a file, just return it if we can handle it
-        if path.is_file():
-            if self.can_handle(path):
-                return [str(path)]
-            return []
-
-        if not path.is_dir():
-            return []
-
-        seen_files: set[str] = set()
-        file_paths: list[str] = []
-
-        for importer in set(self._importers.values()):
-            for file_path in importer.list_dag_files(directory, safe_mode):
-                if file_path not in seen_files:
-                    seen_files.add(file_path)
-                    file_paths.append(file_path)
-
-        return file_paths
-
-    @classmethod
-    def reset(cls) -> None:
-        """Reset the singleton (for testing)."""
-        cls._instance = None
-
-
-def get_importer_registry() -> DagImporterRegistry:
-    """Get the global importer registry instance."""
-    return DagImporterRegistry()
diff --git 
a/airflow-core/src/airflow/dag_processing/importers/python_importer.py 
b/airflow-core/src/airflow/dag_processing/importers/python_importer.py
deleted file mode 100644
index 2eedc4dff4d..00000000000
--- a/airflow-core/src/airflow/dag_processing/importers/python_importer.py
+++ /dev/null
@@ -1,360 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements.  See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership.  The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License.  You may obtain a copy of the License at
-#
-#   http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied.  See the License for the
-# specific language governing permissions and limitations
-# under the License.
-"""Python DAG importer - imports DAGs from Python files."""
-
-from __future__ import annotations
-
-import contextlib
-import importlib
-import importlib.machinery
-import importlib.util
-import logging
-import os
-import signal
-import sys
-import traceback
-import warnings
-import zipfile
-from collections.abc import Iterator
-from pathlib import Path
-from typing import TYPE_CHECKING, Any
-
-from airflow import settings
-from airflow._shared.module_loading.file_discovery import 
find_path_from_directory
-from airflow.configuration import conf
-from airflow.dag_processing.importers.base import (
-    AbstractDagImporter,
-    DagImportError,
-    DagImportResult,
-    DagImportWarning,
-)
-from airflow.utils.docs import get_docs_url
-from airflow.utils.file import get_unique_dag_module_name, might_contain_dag
-
-if TYPE_CHECKING:
-    from types import ModuleType
-
-    from airflow.sdk import DAG
-
-log = logging.getLogger(__name__)
-
-
[email protected]
-def _timeout(seconds: float = 1, error_message: str = "Timeout"):
-    """Context manager for timing out operations."""
-    error_message = error_message + ", PID: " + str(os.getpid())
-
-    def handle_timeout(signum, frame):
-        log.error("Process timed out, PID: %s", str(os.getpid()))
-        from airflow.sdk.exceptions import AirflowTaskTimeout
-
-        raise AirflowTaskTimeout(error_message)
-
-    try:
-        try:
-            signal.signal(signal.SIGALRM, handle_timeout)
-            signal.setitimer(signal.ITIMER_REAL, seconds)
-        except ValueError:
-            log.warning("timeout can't be used in the current context", 
exc_info=True)
-        yield
-    finally:
-        with contextlib.suppress(ValueError):
-            signal.setitimer(signal.ITIMER_REAL, 0)
-
-
-class PythonDagImporter(AbstractDagImporter):
-    """
-    Importer for Python DAG files and zip archives containing Python DAGs.
-
-    This is the default importer registered with the DagImporterRegistry. It 
handles:
-    - .py files: Standard Python DAG files
-    - .zip files: ZIP archives containing Python DAG files
-
-    Note: The .zip extension is exclusively owned by this importer. If you 
need to
-    support other file formats inside ZIP archives (e.g., YAML), you would 
need to
-    either extend this importer or create a composite importer that delegates 
based
-    on the contents of the archive.
-    """
-
-    @classmethod
-    def supported_extensions(cls) -> list[str]:
-        """Return file extensions handled by this importer (.py and .zip)."""
-        return [".py", ".zip"]
-
-    def list_dag_files(
-        self,
-        directory: str | os.PathLike[str],
-        safe_mode: bool = True,
-    ) -> Iterator[str]:
-        """
-        List Python DAG files in a directory.
-
-        Handles both .py files and .zip archives containing Python DAGs.
-        Respects .airflowignore files in the directory tree.
-        """
-        ignore_file_syntax = conf.get_mandatory_value("core", 
"DAG_IGNORE_FILE_SYNTAX", fallback="glob")
-
-        for file_path in find_path_from_directory(directory, ".airflowignore", 
ignore_file_syntax):
-            path = Path(file_path)
-            try:
-                if path.is_file() and (path.suffix.lower() == ".py" or 
zipfile.is_zipfile(path)):
-                    if might_contain_dag(file_path, safe_mode, conf=conf):
-                        yield file_path
-            except Exception:
-                log.exception("Error while examining %s", file_path)
-
-    def import_file(
-        self,
-        file_path: str | Path,
-        *,
-        bundle_path: Path | None = None,
-        bundle_name: str | None = None,
-        safe_mode: bool = True,
-    ) -> DagImportResult:
-        """
-        Import DAGs from a Python file or zip archive.
-
-        :param file_path: Path to the Python file to import.
-        :param bundle_path: Path to the bundle root.
-        :param bundle_name: Name of the bundle.
-        :param safe_mode: If True, skip files that don't appear to contain 
DAGs.
-        :return: DagImportResult with imported DAGs and any errors.
-        """
-        from airflow.sdk.definitions._internal.contextmanager import DagContext
-
-        filepath = str(file_path)
-        relative_path = self.get_relative_path(filepath, bundle_path)
-        result = DagImportResult(file_path=relative_path)
-
-        if not os.path.isfile(filepath):
-            result.errors.append(
-                DagImportError(
-                    file_path=relative_path,
-                    message=f"File not found: {filepath}",
-                    error_type="file_not_found",
-                )
-            )
-            return result
-
-        # Clear any autoregistered dags from previous imports
-        DagContext.autoregistered_dags.clear()
-
-        # Capture warnings during import
-        captured_warnings: list[warnings.WarningMessage] = []
-
-        try:
-            with warnings.catch_warnings(record=True) as captured_warnings:
-                if filepath.endswith(".py") or not 
zipfile.is_zipfile(filepath):
-                    modules = self._load_modules_from_file(filepath, 
safe_mode, result)
-                else:
-                    modules = self._load_modules_from_zip(filepath, safe_mode, 
result)
-        except TypeError:
-            # Configuration errors (e.g., invalid timeout type) should 
propagate
-            raise
-        except Exception as e:
-            result.errors.append(
-                DagImportError(
-                    file_path=relative_path,
-                    message=str(e),
-                    error_type="import",
-                    stacktrace=traceback.format_exc(),
-                )
-            )
-            return result
-
-        # Convert captured warnings to DagImportWarning
-        for warn_msg in captured_warnings:
-            category = warn_msg.category.__name__
-            if (module := warn_msg.category.__module__) != "builtins":
-                category = f"{module}.{category}"
-            result.warnings.append(
-                DagImportWarning(
-                    file_path=warn_msg.filename,
-                    message=str(warn_msg.message),
-                    warning_type=category,
-                    line_number=warn_msg.lineno,
-                )
-            )
-
-        # Process imported modules to extract DAGs
-        self._process_modules(filepath, modules, bundle_name, bundle_path, 
result)
-
-        return result
-
-    def _load_modules_from_file(
-        self, filepath: str, safe_mode: bool, result: DagImportResult
-    ) -> list[ModuleType]:
-        from airflow.sdk.definitions._internal.contextmanager import DagContext
-
-        def sigsegv_handler(signum, frame):
-            msg = f"Received SIGSEGV signal while processing {filepath}."
-            log.error(msg)
-            result.errors.append(
-                DagImportError(
-                    file_path=result.file_path,
-                    message=msg,
-                    error_type="segfault",
-                )
-            )
-
-        try:
-            signal.signal(signal.SIGSEGV, sigsegv_handler)
-        except ValueError:
-            log.warning("SIGSEGV signal handler registration failed. Not in 
the main thread")
-
-        if not might_contain_dag(filepath, safe_mode, conf=conf):
-            log.debug("File %s assumed to contain no DAGs. Skipping.", 
filepath)
-            result.skipped_files.append(filepath)
-            return []
-
-        log.debug("Importing %s", filepath)
-        mod_name = get_unique_dag_module_name(filepath)
-
-        if mod_name in sys.modules:
-            del sys.modules[mod_name]
-
-        DagContext.current_autoregister_module_name = mod_name
-
-        def parse(mod_name: str, filepath: str) -> list[ModuleType]:
-            try:
-                loader = importlib.machinery.SourceFileLoader(mod_name, 
filepath)
-                spec = importlib.util.spec_from_loader(mod_name, loader)
-                new_module = importlib.util.module_from_spec(spec)  # type: 
ignore[arg-type]
-                sys.modules[spec.name] = new_module  # type: ignore[union-attr]
-                loader.exec_module(new_module)
-                return [new_module]
-            except KeyboardInterrupt:
-                raise
-            except BaseException as e:
-                DagContext.autoregistered_dags.clear()
-                log.exception("Failed to import: %s", filepath)
-                if conf.getboolean("core", "dagbag_import_error_tracebacks"):
-                    stacktrace = traceback.format_exc(
-                        limit=-conf.getint("core", 
"dagbag_import_error_traceback_depth")
-                    )
-                else:
-                    stacktrace = None
-                result.errors.append(
-                    DagImportError(
-                        file_path=result.file_path,
-                        message=str(e),
-                        error_type="import",
-                        stacktrace=stacktrace,
-                    )
-                )
-                return []
-
-        dagbag_import_timeout = settings.get_dagbag_import_timeout(filepath)
-
-        if not isinstance(dagbag_import_timeout, (int, float)):
-            raise TypeError(
-                f"Value ({dagbag_import_timeout}) from 
get_dagbag_import_timeout must be int or float"
-            )
-
-        if dagbag_import_timeout <= 0:
-            return parse(mod_name, filepath)
-
-        timeout_msg = (
-            f"DagBag import timeout for {filepath} after 
{dagbag_import_timeout}s.\n"
-            "Please take a look at these docs to improve your DAG import 
time:\n"
-            f"* {get_docs_url('best-practices.html#top-level-python-code')}\n"
-            f"* {get_docs_url('best-practices.html#reducing-dag-complexity')}"
-        )
-        with _timeout(dagbag_import_timeout, error_message=timeout_msg):
-            return parse(mod_name, filepath)
-
-    def _load_modules_from_zip(
-        self, filepath: str, safe_mode: bool, result: DagImportResult
-    ) -> list[ModuleType]:
-        """Load Python modules from a zip archive."""
-        from airflow.sdk.definitions._internal.contextmanager import DagContext
-
-        mods: list[ModuleType] = []
-        with zipfile.ZipFile(filepath) as current_zip_file:
-            for zip_info in current_zip_file.infolist():
-                zip_path = Path(zip_info.filename)
-                if zip_path.suffix not in [".py", ".pyc"] or 
len(zip_path.parts) > 1:
-                    continue
-
-                if zip_path.stem == "__init__":
-                    log.warning("Found %s at root of %s", zip_path.name, 
filepath)
-
-                log.debug("Reading %s from %s", zip_info.filename, filepath)
-
-                if not might_contain_dag(zip_info.filename, safe_mode, 
current_zip_file, conf=conf):
-                    
result.skipped_files.append(f"{filepath}:{zip_info.filename}")
-                    continue
-
-                mod_name = zip_path.stem
-                if mod_name in sys.modules:
-                    del sys.modules[mod_name]
-
-                DagContext.current_autoregister_module_name = mod_name
-                try:
-                    sys.path.insert(0, filepath)
-                    current_module = importlib.import_module(mod_name)
-                    mods.append(current_module)
-                except Exception as e:
-                    DagContext.autoregistered_dags.clear()
-                    fileloc = os.path.join(filepath, zip_info.filename)
-                    log.exception("Failed to import: %s", fileloc)
-                    if conf.getboolean("core", 
"dagbag_import_error_tracebacks"):
-                        stacktrace = traceback.format_exc(
-                            limit=-conf.getint("core", 
"dagbag_import_error_traceback_depth")
-                        )
-                    else:
-                        stacktrace = None
-                    result.errors.append(
-                        DagImportError(
-                            file_path=fileloc,  # Use the file path inside the 
ZIP
-                            message=str(e),
-                            error_type="import",
-                            stacktrace=stacktrace,
-                        )
-                    )
-                finally:
-                    if sys.path[0] == filepath:
-                        del sys.path[0]
-        return mods
-
-    def _process_modules(
-        self,
-        filepath: str,
-        mods: list[Any],
-        bundle_name: str | None,
-        bundle_path: Path | None,
-        result: DagImportResult,
-    ) -> None:
-        """Extract DAG objects from modules. Validation happens in 
bag_dag()."""
-        from airflow.sdk import DAG
-        from airflow.sdk.definitions._internal.contextmanager import DagContext
-
-        top_level_dags: set[tuple[DAG, Any]] = {
-            (o, m) for m in mods for o in m.__dict__.values() if isinstance(o, 
DAG)
-        }
-        top_level_dags.update(DagContext.autoregistered_dags)
-
-        DagContext.current_autoregister_module_name = None
-        DagContext.autoregistered_dags.clear()
-
-        for dag, mod in top_level_dags:
-            dag.fileloc = mod.__file__
-            relative_fileloc = self.get_relative_path(dag.fileloc, bundle_path)
-            dag.relative_fileloc = relative_fileloc
-
-            result.dags.append(dag)
-            log.debug("Found DAG %s", dag.dag_id)
diff --git a/airflow-core/src/airflow/dag_processing/manager.py 
b/airflow-core/src/airflow/dag_processing/manager.py
index c6ce3c2dcbe..8c4fd0588eb 100644
--- a/airflow-core/src/airflow/dag_processing/manager.py
+++ b/airflow-core/src/airflow/dag_processing/manager.py
@@ -1355,6 +1355,7 @@ class DagFileProcessorManager(LoggingMixin):
         files_parsed: set[tuple[str, str]] | None = None
         if relative_fileloc is not None:
             files_parsed = {(bundle_name, relative_fileloc)}
+            files_parsed.update((bundle_name, rel_path) for rel_path in 
parsing_result.parsed_definitions)
             files_parsed.update(import_errors.keys())
 
         warnings = parsing_result.warnings or []
@@ -1371,6 +1372,7 @@ class DagFileProcessorManager(LoggingMixin):
             warnings=set(warnings),
             session=session,
             files_parsed=files_parsed,
+            dag_source_codes=parsing_result.dag_source_codes,
         )
 
     def _collect_results(self):
diff --git a/airflow-core/src/airflow/dag_processing/processor.py 
b/airflow-core/src/airflow/dag_processing/processor.py
index fe29a314100..80e6d7161de 100644
--- a/airflow-core/src/airflow/dag_processing/processor.py
+++ b/airflow-core/src/airflow/dag_processing/processor.py
@@ -73,6 +73,7 @@ from airflow.sdk.execution_time.comms import (
 )
 from airflow.sdk.execution_time.supervisor import WatchedSubprocess, 
register_request_method
 from airflow.sdk.execution_time.task_runner import RuntimeTaskInstance, 
_send_error_email_notification
+from airflow.sdk.importers import DagSourceCode  # noqa: TC001
 from airflow.serialization.serialized_objects import DagSerialization, 
LazyDeserializedDAG
 from airflow.utils.dag_version_inflation_checker import 
check_dag_file_stability
 from airflow.utils.file import iter_airflow_imports
@@ -127,6 +128,10 @@ class DagFileParsingResult(BaseModel):
     serialized_dags: list[LazyDeserializedDAG]
     warnings: list | None = None
     import_errors: dict[str, str] | None = None
+    parsed_definitions: list[str] = Field(default_factory=list)
+    """Bundle-relative locations of the Dag definitions imported from 
``fileloc``."""
+    dag_source_codes: dict[str, DagSourceCode] = Field(default_factory=dict)
+    """Source code of the parsed Dags, keyed by Dag fileloc."""
     type: Literal["DagFileParsingResult"] = "DagFileParsingResult"
 
 
@@ -253,7 +258,15 @@ def _parse_file(msg: DagFileParseRequest, log: 
FilteringBoundLogger) -> DagFileP
         fileloc=msg.file,
         serialized_dags=serialized_dags,
         import_errors=bag.import_errors,
-        warnings=stability_check_result.get_formatted_warnings(bag.dag_ids),
+        warnings=[
+            *stability_check_result.get_formatted_warnings(bag.dag_ids),
+            *(
+                {"dag_id": w.dag_id, "warning_type": w.warning_type, 
"message": w.message}
+                for w in bag.dag_warnings
+            ),
+        ],
+        parsed_definitions=bag.parsed_definitions,
+        dag_source_codes=bag.dag_source_codes,
     )
     return result
 
diff --git a/airflow-core/src/airflow/models/dagcode.py 
b/airflow-core/src/airflow/models/dagcode.py
index b21af49d57f..12c02ed00f1 100644
--- a/airflow-core/src/airflow/models/dagcode.py
+++ b/airflow-core/src/airflow/models/dagcode.py
@@ -42,6 +42,7 @@ if TYPE_CHECKING:
     from sqlalchemy.sql import Select
 
     from airflow.models.dag_version import DagVersion
+    from airflow.sdk.importers import DagSourceCode  # noqa: SDK001
 
 log = logging.getLogger(__name__)
 
@@ -75,24 +76,40 @@ class DagCode(Base):
     dag_version = relationship("DagVersion", back_populates="dag_code", 
uselist=False)
     __table_args__ = (Index("idx_dag_code_dag_id_last_updated", dag_id, 
last_updated),)
 
-    def __init__(self, dag_version, full_filepath: str, source_code: str | 
None = None):
+    def __init__(
+        self,
+        dag_version: DagVersion,
+        full_filepath: str,
+        source_code: str | None = None,
+        language: str = "python",
+    ):
         self.dag_version = dag_version
         self.fileloc = full_filepath
         self.source_code = source_code or DagCode.code(self.dag_version.dag_id)
         self.source_code_hash = self.dag_source_hash(self.source_code)
         self.dag_id = dag_version.dag_id
+        self.language = language
 
     @classmethod
     @provide_session
-    def write_code(cls, dag_version: DagVersion, fileloc: str, *, session: 
Session = NEW_SESSION) -> DagCode:
+    def write_code(
+        cls,
+        dag_version: DagVersion,
+        fileloc: str,
+        *,
+        dag_source_code: DagSourceCode | None = None,
+        session: Session = NEW_SESSION,
+    ) -> DagCode:
         """
         Write code into database.
 
         :param fileloc: file path of DAG to sync
+        :param dag_source_code: Source code read by the Dag importer; read 
from ``fileloc`` when not given
         :param session: ORM Session
         """
         log.debug("Writing DAG file %s into DagCode table", fileloc)
-        dag_code = DagCode(dag_version, fileloc, 
cls.get_code_from_file(fileloc))
+        source_code, language = cls._get_source_code_and_language(fileloc, 
dag_source_code)
+        dag_code = DagCode(dag_version, fileloc, source_code, 
language=language)
         session.add(dag_code)
         log.debug("DAG file %s written into DagCode table", fileloc)
         return dag_code
@@ -133,6 +150,14 @@ class DagCode(Base):
                 return "source_code"
             raise
 
+    @classmethod
+    def _get_source_code_and_language(
+        cls, fileloc: str, dag_source_code: DagSourceCode | None
+    ) -> tuple[str, str]:
+        if dag_source_code is None:
+            return cls.get_code_from_file(fileloc), "python"
+        return dag_source_code.source_code, dag_source_code.language
+
     @classmethod
     @provide_session
     def _get_code_from_db(cls, dag_id, *, session: Session = NEW_SESSION) -> 
str:
@@ -177,23 +202,33 @@ class DagCode(Base):
 
     @classmethod
     @provide_session
-    def update_source_code(cls, dag_id: str, fileloc: str, *, session: Session 
= NEW_SESSION) -> None:
+    def update_source_code(
+        cls,
+        dag_id: str,
+        fileloc: str,
+        *,
+        dag_source_code: DagSourceCode | None = None,
+        session: Session = NEW_SESSION,
+    ) -> None:
         """
         Check if the source code of the DAG has changed and update it if 
needed.
 
         :param dag_id: Dag ID
         :param fileloc: The path of code file to read the code from
+        :param dag_source_code: Source code read by the Dag importer; read 
from ``fileloc`` when not given
         :param session: The database session.
         :return: None
         """
         latest_dagcode = cls.get_latest_dagcode(dag_id, session=session)
         if not latest_dagcode:
             return
-        new_source_code = cls.get_code_from_file(fileloc)
+        new_source_code, new_language = 
cls._get_source_code_and_language(fileloc, dag_source_code)
         new_source_code_hash = cls.dag_source_hash(new_source_code)
         if new_source_code_hash != latest_dagcode.source_code_hash:
             latest_dagcode.source_code = new_source_code
             latest_dagcode.source_code_hash = new_source_code_hash
+        if latest_dagcode.language != new_language:
+            latest_dagcode.language = new_language
         # Keep fileloc aligned even when the contents are unchanged (e.g. the 
file was moved/renamed).
         if fileloc != latest_dagcode.fileloc:
             latest_dagcode.fileloc = fileloc
diff --git a/airflow-core/src/airflow/models/serialized_dag.py 
b/airflow-core/src/airflow/models/serialized_dag.py
index d49c168013f..b933194470d 100644
--- a/airflow-core/src/airflow/models/serialized_dag.py
+++ b/airflow-core/src/airflow/models/serialized_dag.py
@@ -63,6 +63,7 @@ if TYPE_CHECKING:
     from sqlalchemy.sql import Select
     from sqlalchemy.sql.elements import ColumnElement
 
+    from airflow.sdk.importers import DagSourceCode  # noqa: SDK001
     from airflow.serialization.definitions.dag import SerializedDAG
     from airflow.serialization.serialized_objects import LazyDeserializedDAG
 
@@ -612,6 +613,7 @@ class SerializedDagModel(Base):
         version_data: dict | None = None,
         min_update_interval: int | None = None,
         *,
+        dag_source_code: DagSourceCode | None = None,
         session: Session = NEW_SESSION,
         _prefetched: DagWriteMetadata | None = None,
     ) -> bool:
@@ -626,6 +628,7 @@ class SerializedDagModel(Base):
         :param bundle_version: bundle version of the DAG
         :param version_data: optional structured data associated with this 
version
         :param min_update_interval: minimal interval in seconds to update 
serialized DAG
+        :param dag_source_code: Source code read by the Dag importer; read 
from ``fileloc`` when not given
         :param session: ORM Session
         :param _prefetched: Pre-fetched metadata to skip per-DAG queries; used 
by bulk callers
 
@@ -713,7 +716,12 @@ class SerializedDagModel(Base):
                 dag_version.bundle_version = bundle_version
                 dag_version.version_data = version_data
                 session.merge(dag_version)
-                DagCode.update_source_code(dag_id=dag.dag_id, 
fileloc=dag.fileloc, session=session)
+                DagCode.update_source_code(
+                    dag_id=dag.dag_id,
+                    fileloc=dag.fileloc,
+                    dag_source_code=dag_source_code,
+                    session=session,
+                )
             if name_updated or bundle_metadata_changed:
                 # A write occurred — a deadline alert name update and/or a 
bundle
                 # metadata refresh — so report True so callers know the DB 
changed.
@@ -774,7 +782,12 @@ class SerializedDagModel(Base):
             dag_version.version_data = version_data
             session.merge(dag_version)
             # Update the latest DagCode
-            DagCode.update_source_code(dag_id=dag.dag_id, fileloc=dag.fileloc, 
session=session)
+            DagCode.update_source_code(
+                dag_id=dag.dag_id,
+                fileloc=dag.fileloc,
+                dag_source_code=dag_source_code,
+                session=session,
+            )
             stats.incr(
                 "dag.serialization.version_updated",
                 tags={"dag_id": dag.dag_id, "bundle_name": bundle_name},
@@ -802,7 +815,7 @@ class SerializedDagModel(Base):
 
         cls._create_deadline_alert_records(new_serialized_dag, 
deadline_uuid_mapping)
         log.debug("DAG: %s written to the DB", dag.dag_id)
-        DagCode.write_code(dagv, dag.fileloc, session=session)
+        DagCode.write_code(dagv, dag.fileloc, dag_source_code=dag_source_code, 
session=session)
         stats.incr(
             "dag.serialization.version_created",
             tags={"dag_id": dag.dag_id, "bundle_name": bundle_name},
diff --git a/airflow-core/src/airflow/utils/file.py 
b/airflow-core/src/airflow/utils/file.py
index c762fffaf9f..ff7d7db1662 100644
--- a/airflow-core/src/airflow/utils/file.py
+++ b/airflow-core/src/airflow/utils/file.py
@@ -119,6 +119,16 @@ def find_dag_file_paths(directory: str | os.PathLike[str], 
safe_mode: bool) -> l
     return file_paths
 
 
+def find_enclosing_file(path: Path) -> Path | None:
+    """
+    Return ``path`` or its nearest ancestor that is a file, or ``None`` if 
there is none.
+
+    A Dag definition nested in a container is referenced by a path that does 
not exist on
+    disk (``archive.zip/dag.py``, for instance); this resolves it to the 
container.
+    """
+    return next((candidate for candidate in (path, *path.parents) if 
candidate.is_file()), None)
+
+
 COMMENT_PATTERN = re.compile(r"\s*#.*")
 
 
diff --git a/airflow-core/tests/unit/dag_processing/importers/__init__.py 
b/airflow-core/tests/unit/dag_processing/importers/__init__.py
deleted file mode 100644
index 13a83393a91..00000000000
--- a/airflow-core/tests/unit/dag_processing/importers/__init__.py
+++ /dev/null
@@ -1,16 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements.  See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership.  The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License.  You may obtain a copy of the License at
-#
-#   http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied.  See the License for the
-# specific language governing permissions and limitations
-# under the License.
diff --git a/airflow-core/tests/unit/dag_processing/importers/test_registry.py 
b/airflow-core/tests/unit/dag_processing/importers/test_registry.py
deleted file mode 100644
index 955308b4559..00000000000
--- a/airflow-core/tests/unit/dag_processing/importers/test_registry.py
+++ /dev/null
@@ -1,99 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements.  See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership.  The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License.  You may obtain a copy of the License at
-#
-#   http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing,
-# software distributed under the License is distributed on an
-# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-# KIND, either express or implied.  See the License for the
-# specific language governing permissions and limitations
-# under the License.
-"""Tests for the DagImporterRegistry."""
-
-from __future__ import annotations
-
-from pathlib import Path
-
-from airflow.dag_processing.importers import (
-    DagImporterRegistry,
-    PythonDagImporter,
-    get_importer_registry,
-)
-
-
-class TestDagImporterRegistry:
-    """Test the DagImporterRegistry singleton."""
-
-    def setup_method(self):
-        """Reset the registry before each test."""
-        DagImporterRegistry.reset()
-
-    def teardown_method(self):
-        """Reset the registry after each test."""
-        DagImporterRegistry.reset()
-
-    def test_singleton_pattern(self):
-        """Registry should return the same instance."""
-        registry1 = get_importer_registry()
-        registry2 = get_importer_registry()
-        assert registry1 is registry2
-
-    def test_default_importers_registered(self):
-        """Registry should have Python importer by default."""
-        registry = get_importer_registry()
-        extensions = registry.supported_extensions()
-
-        assert ".py" in extensions
-
-    def test_get_importer_for_python(self):
-        """Should return PythonDagImporter for .py files."""
-        registry = get_importer_registry()
-        importer = registry.get_importer("test.py")
-
-        assert importer is not None
-        assert isinstance(importer, PythonDagImporter)
-
-    def test_get_importer_for_unknown(self):
-        """Should return None for unknown file types."""
-        registry = get_importer_registry()
-        importer = registry.get_importer("test.txt")
-
-        assert importer is None
-
-    def test_can_handle_supported_files(self):
-        """can_handle should return True for supported file types."""
-        registry = get_importer_registry()
-
-        assert registry.can_handle("dag.py")
-        assert registry.can_handle(Path("subdir/dag.py"))
-
-    def test_can_handle_unsupported_files(self):
-        """can_handle should return False for unsupported file types."""
-        registry = get_importer_registry()
-
-        assert not registry.can_handle("readme.txt")
-        assert not registry.can_handle("config.json")
-        assert not registry.can_handle("script.sh")
-
-    def test_case_insensitive_extension_matching(self):
-        """Extension matching should be case-insensitive."""
-        registry = get_importer_registry()
-
-        # All these should be handled
-        assert registry.can_handle("dag.PY")
-        assert registry.can_handle("dag.Py")
-
-    def test_reset_clears_singleton(self):
-        """reset() should clear the singleton instance."""
-        registry1 = get_importer_registry()
-        DagImporterRegistry.reset()
-        registry2 = get_importer_registry()
-
-        # Should be different instances after reset
-        assert registry1 is not registry2
diff --git a/airflow-core/tests/unit/dag_processing/test_collection.py 
b/airflow-core/tests/unit/dag_processing/test_collection.py
index da9342958d6..d7e5213fad5 100644
--- a/airflow-core/tests/unit/dag_processing/test_collection.py
+++ b/airflow-core/tests/unit/dag_processing/test_collection.py
@@ -61,6 +61,7 @@ from airflow.models.asset import (
 )
 from airflow.models.dag import DagTag
 from airflow.models.dagbundle import DagBundleModel
+from airflow.models.dagcode import DagCode
 from airflow.models.dagwarning import DagWarning, DagWarningType
 from airflow.models.errors import ParseImportError
 from airflow.models.serialized_dag import SerializedDagModel
@@ -83,6 +84,7 @@ from airflow.sdk import (
 )
 from airflow.sdk.definitions.deadline import AsyncCallback, 
BaseDeadlineReference, DeadlineAlert
 from airflow.sdk.definitions.timetables.assets import AssetOrTimeSchedule, 
PartitionedAssetTimetable
+from airflow.sdk.importers import DagSourceCode
 from airflow.serialization.definitions.assets import SerializedAsset
 from airflow.serialization.encoders import encode_trigger, 
ensure_serialized_asset
 from airflow.serialization.serialized_objects import LazyDeserializedDAG
@@ -800,6 +802,7 @@ class TestUpdateDagParsingResults:
                     bundle_version=None,
                     version_data=None,
                     min_update_interval=mock.ANY,
+                    dag_source_code=None,
                     session=mock_session,
                     _prefetched=mock.ANY,
                 ),
@@ -910,6 +913,55 @@ class TestUpdateDagParsingResults:
 
         assert warning is None
 
+    @pytest.mark.usefixtures("clean_db")
+    def test_stale_importer_warnings_are_replaced(self, testing_dag_bundle, 
session):
+        session.add(DagModel(dag_id="imported_dag", bundle_name="testing", 
fileloc="/dags/imported.py"))
+        session.flush()
+        session.add_all(
+            [
+                DagWarning(dag_id="imported_dag", warning_type="test:stale", 
message="Stale"),
+                DagWarning(
+                    dag_id="imported_dag", 
warning_type=DagWarningType.ASSET_CONFLICT, message="Conflict"
+                ),
+            ]
+        )
+        session.flush()
+
+        update_dag_parsing_results_in_db(
+            bundle_name="testing",
+            bundle_version=None,
+            dags=[LazyDeserializedDAG.from_dag(DAG(dag_id="imported_dag"))],
+            import_errors={},
+            parse_duration=None,
+            warnings={DagWarning("imported_dag", "test:current", "Current")},
+            session=session,
+        )
+
+        warning_types = session.scalars(
+            select(DagWarning.warning_type).where(DagWarning.dag_id == 
"imported_dag")
+        ).all()
+        assert sorted(warning_types) == [DagWarningType.ASSET_CONFLICT.value, 
"test:current"]
+
+    @pytest.mark.usefixtures("clean_db")
+    def test_dag_source_codes_are_written_to_dag_code(self, 
testing_dag_bundle, session):
+        dag = DAG(dag_id="yaml_dag")
+        dag.fileloc = "/dags/yaml_dag.yaml"
+        dag.relative_fileloc = "yaml_dag.yaml"
+
+        update_dag_parsing_results_in_db(
+            bundle_name="testing",
+            bundle_version=None,
+            dags=[LazyDeserializedDAG.from_dag(dag)],
+            import_errors={},
+            parse_duration=None,
+            warnings=set(),
+            session=session,
+            dag_source_codes={dag.fileloc: DagSourceCode(source_code="dag_id: 
yaml_dag\n", language="yaml")},
+        )
+
+        dag_code = DagCode.get_latest_dagcode("yaml_dag", session=session)
+        assert (dag_code.source_code, dag_code.language) == ("dag_id: 
yaml_dag\n", "yaml")
+
     def test_parse_time_written_to_db_on_sync(self, testing_dag_bundle, 
session):
         """Test that the parse time is correctly written to the DB after 
parsing"""
 
diff --git a/airflow-core/tests/unit/dag_processing/test_dagbag.py 
b/airflow-core/tests/unit/dag_processing/test_dagbag.py
index cb110ef5c7f..54b06ce75aa 100644
--- a/airflow-core/tests/unit/dag_processing/test_dagbag.py
+++ b/airflow-core/tests/unit/dag_processing/test_dagbag.py
@@ -18,7 +18,7 @@ from __future__ import annotations
 
 import contextlib
 import inspect
-import logging
+import json
 import os
 import pathlib
 import re
@@ -39,8 +39,8 @@ from airflow import settings
 from airflow.dag_processing.dagbag import (
     BundleDagBag,
     DagBag,
-    _capture_with_reraise,
     _validate_executor_fields,
+    sync_bag_to_db,
 )
 from airflow.exceptions import UnknownExecutorException
 from airflow.executors.executor_loader import ExecutorLoader
@@ -49,6 +49,16 @@ from airflow.models.dagwarning import DagWarning, 
DagWarningType
 from airflow.models.pool import Pool
 from airflow.models.serialized_dag import SerializedDagModel
 from airflow.sdk import DAG, BaseOperator
+from airflow.sdk.exceptions import AirflowConfigException
+from airflow.sdk.importers import (
+    AbstractDagImporter,
+    DagImportResult,
+    DagImportWarning,
+    DagSourceCode,
+    PythonDagImporter,
+    find_file_dag_definitions,
+    get_file_suffix,
+)
 
 from tests_common.pytest_plugin import AIRFLOW_ROOT_PATH
 from tests_common.test_utils import db
@@ -345,6 +355,39 @@ def test_validate_executor_field():
         _validate_executor_fields(dag)
 
 
+def _dag_source(dag_id: str) -> str:
+    return f"from airflow.sdk import DAG\n\ndag = DAG({dag_id!r}, 
schedule=None)\n"
+
+
+class TextDagImporter(AbstractDagImporter):
+    """Builds one Dag per ``.dagtxt`` file, whose content is the dag_id."""
+
+    supported_extensions = [".dagtxt"]
+
+    def can_handle(self, definition):
+        return get_file_suffix(definition) == ".dagtxt"
+
+    def list_dag_definitions(self, bundle, *, safe_mode=True):
+        yield from find_file_dag_definitions(bundle.path, 
self.supported_extensions)
+
+    def import_definition(self, definition, bundle):
+        return DagImportResult(
+            definition=definition,
+            dags=[DAG(definition.read_text().strip(), schedule=None)],
+            warnings=[
+                DagImportWarning(
+                    source_reference=repr(definition),
+                    message="Deprecated field",
+                    warning_type="test:deprecated_field",
+                    line_number=1,
+                )
+            ],
+        )
+
+    def get_source_code(self, definition):
+        return DagSourceCode(source_code=definition.read_text(), 
language="text")
+
+
 class TestDagBag:
     def setup_class(self):
         db_clean_up()
@@ -686,19 +729,77 @@ class TestDagBag:
         assert os.fspath(path2) not in dagbag.import_errors
         assert "AirflowDagDuplicatedIdException" in 
dagbag.import_errors[error_path]
 
-    def test_zip_skip_log(self, caplog, test_zip_path):
-        """
-        test the loading of a DAG from within a zip file that skips another 
file because
-        it doesn't have "airflow" and "DAG"
-        """
-        caplog.set_level(logging.INFO)
+    def test_zip_skips_members_without_dag_markers(self, test_zip_path):
         dagbag = DagBag(dag_folder=test_zip_path)
 
-        assert dagbag.has_logged
-        assert (
-            f"File {test_zip_path}:file_no_airflow_dag.py "
-            "assumed to contain no DAGs. Skipping." in caplog.text
-        )
+        assert f"{test_zip_path}/test_zip.py" in dagbag.parsed_definitions
+        assert f"{test_zip_path}/file_no_airflow_dag.py" not in 
dagbag.parsed_definitions
+
+    def test_zip_member_path_imports_only_that_member(self, tmp_path):
+        zipped = tmp_path / "dags.zip"
+        with zipfile.ZipFile(zipped, "w") as zf:
+            zf.writestr("first.py", _dag_source("first"))
+            zf.writestr("second.py", _dag_source("second"))
+
+        dagbag = DagBag(dag_folder=zipped / "first.py", bundle_path=tmp_path, 
bundle_name="test_bundle")
+
+        assert dagbag.dag_ids == ["first"]
+        assert dagbag.dags["first"].relative_fileloc == "dags.zip/first.py"
+        assert dagbag.parsed_definitions == ["dags.zip/first.py"]
+
+    @pytest.mark.parametrize("dag_folder", ["", "corrupt.zip"], ids=["bundle", 
"archive"])
+    def test_discovery_errors_use_relative_path_with_bundle(self, tmp_path, 
dag_folder):
+        (tmp_path / "corrupt.zip").write_bytes(b"not a zip")
+
+        dagbag = DagBag(dag_folder=tmp_path / dag_folder, 
bundle_path=tmp_path, bundle_name="test_bundle")
+
+        assert list(dagbag.import_errors) == ["corrupt.zip"]
+
+    def test_dag_source_codes(self, tmp_path):
+        (tmp_path / "plain.py").write_text(_dag_source("plain"))
+        with zipfile.ZipFile(tmp_path / "dags.zip", "w") as zf:
+            zf.writestr("member.py", _dag_source("member"))
+
+        dagbag = DagBag(dag_folder=tmp_path, bundle_path=tmp_path, 
bundle_name="test_bundle")
+
+        assert dagbag.dag_source_codes == {
+            os.fspath(tmp_path / "plain.py"): DagSourceCode(
+                source_code=_dag_source("plain"), language="python"
+            ),
+            os.fspath(tmp_path / "dags.zip" / "member.py"): DagSourceCode(
+                source_code=_dag_source("member"), language="python"
+            ),
+        }
+
+    def test_dag_source_codes_drops_stale_source_when_read_fails(self, 
tmp_path):
+        dag_file = tmp_path / "plain.py"
+        dag_file.write_text(_dag_source("plain"))
+        dagbag = DagBag(dag_folder=tmp_path, bundle_path=tmp_path, 
bundle_name="test_bundle")
+        dag_file.write_text(_dag_source("plain") + "# changed\n")
+
+        with mock.patch.object(
+            PythonDagImporter, "get_source_code", autospec=True, 
side_effect=OSError("unreadable")
+        ):
+            dagbag.process_file(os.fspath(dag_file), only_if_updated=True)
+
+        assert "plain" in dagbag.dags
+        assert dagbag.dag_source_codes == {}
+
+    
@mock.patch("airflow.dag_processing.collection.update_dag_parsing_results_in_db",
 autospec=True)
+    def test_sync_bag_to_db_passes_parsed_files_and_source_codes(self, 
mock_update, tmp_path):
+        with zipfile.ZipFile(tmp_path / "dags.zip", "w") as zf:
+            zf.writestr("member.py", _dag_source("member"))
+        dagbag = DagBag(dag_folder=tmp_path, bundle_path=tmp_path, 
bundle_name="test_bundle")
+
+        sync_bag_to_db(dagbag, "test_bundle", None, 
session=mock.sentinel.session)
+
+        kwargs = mock_update.call_args.kwargs
+        assert kwargs["files_parsed"] == {("test_bundle", 
"dags.zip/member.py"), ("test_bundle", "dags.zip")}
+        assert kwargs["dag_source_codes"] == {
+            os.fspath(tmp_path / "dags.zip" / "member.py"): DagSourceCode(
+                source_code=_dag_source("member"), language="python"
+            )
+        }
 
     def test_zip(self, tmp_path, test_zip_path):
         """
@@ -711,7 +812,7 @@ class TestDagBag:
         assert sys.path == syspath_before  # sys.path doesn't change
         assert not dagbag.import_errors
 
-    @patch("airflow.dag_processing.importers.python_importer._timeout")
+    @patch("airflow.sdk.importers.python_importer.timeout")
     @patch("airflow.dag_processing.dagbag.settings.get_dagbag_import_timeout")
     def test_process_dag_file_without_timeout(
         self, mocked_get_dagbag_import_timeout, mocked_timeout, tmp_path
@@ -730,7 +831,7 @@ class TestDagBag:
         dagbag.process_file(os.path.join(TEST_DAGS_FOLDER, "test_sensor.py"))
         mocked_timeout.assert_not_called()
 
-    @patch("airflow.dag_processing.importers.python_importer._timeout")
+    @patch("airflow.sdk.importers.python_importer.timeout")
     @patch("airflow.dag_processing.dagbag.settings.get_dagbag_import_timeout")
     def test_process_dag_file_with_non_default_timeout(
         self, mocked_get_dagbag_import_timeout, mocked_timeout, tmp_path
@@ -747,9 +848,9 @@ class TestDagBag:
         dagbag = DagBag(dag_folder=os.fspath(tmp_path))
         dagbag.process_file(os.path.join(TEST_DAGS_FOLDER, "test_sensor.py"))
 
-        mocked_timeout.assert_called_once_with(timeout_value, 
error_message=mock.ANY)
+        mocked_timeout.assert_called_once_with(seconds=timeout_value, 
error_message=mock.ANY)
 
-    
@patch("airflow.dag_processing.importers.python_importer.settings.get_dagbag_import_timeout")
+    @patch("airflow.dag_processing.dagbag.settings.get_dagbag_import_timeout")
     def test_check_value_type_from_get_dagbag_import_timeout(
         self, mocked_get_dagbag_import_timeout, tmp_path
     ):
@@ -760,7 +861,7 @@ class TestDagBag:
 
         dagbag = DagBag(dag_folder=os.fspath(tmp_path))
         with pytest.raises(
-            TypeError, match=r"Value \(1\) from get_dagbag_import_timeout must 
be int or float"
+            AirflowConfigException, match=r"Value \(1\) from 
get_dagbag_import_timeout must be int or float"
         ):
             dagbag.process_file(os.path.join(TEST_DAGS_FOLDER, 
"test_sensor.py"))
 
@@ -833,20 +934,19 @@ class TestDagBag:
         mock_dagmodel.return_value.fileloc = "foo"
 
         class _TestDagBag(DagBag):
-            process_file_calls = 0
+            import_calls = 0
 
-            def process_file(self, filepath, only_if_updated=True, 
safe_mode=True):
-                if os.path.basename(filepath) == "example_bash_operator.py":
-                    _TestDagBag.process_file_calls += 1
-                super().process_file(filepath, only_if_updated, safe_mode)
+            def _process_definition(self, importer, definition, *, 
only_if_updated):
+                if os.path.basename(repr(definition)) == 
"example_bash_operator.py":
+                    _TestDagBag.import_calls += 1
+                return super()._process_definition(importer, definition, 
only_if_updated=only_if_updated)
 
         dagbag = _TestDagBag(dag_folder=standard_example_dags_folder)
-        dagbag.process_file_calls
 
-        # Should not call process_file again, since it's already loaded during 
init.
-        assert dagbag.process_file_calls == 1
+        # Should not import the file again, since it's already loaded during 
init.
+        assert dagbag.import_calls == 1
         assert dagbag.get_dag(dag_id) is not None
-        assert dagbag.process_file_calls == 1
+        assert dagbag.import_calls == 1
 
     @pytest.mark.parametrize(
         ("file_name", "expected_dag_id"),
@@ -922,20 +1022,20 @@ class TestDagBag:
         mock_dagmodel.return_value.fileloc = fileloc
 
         class _TestDagBag(DagBag):
-            process_file_calls = 0
+            import_calls = 0
 
-            def process_file(self, filepath, only_if_updated=True, 
safe_mode=True):
-                if filepath == fileloc:
-                    _TestDagBag.process_file_calls += 1
-                return super().process_file(filepath, only_if_updated, 
safe_mode)
+            def _process_definition(self, importer, definition, *, 
only_if_updated):
+                if repr(definition) == fileloc:
+                    _TestDagBag.import_calls += 1
+                return super()._process_definition(importer, definition, 
only_if_updated=only_if_updated)
 
         dagbag = _TestDagBag(dag_folder=standard_example_dags_folder)
 
-        assert dagbag.process_file_calls == 1
+        assert dagbag.import_calls == 1
         dag = dagbag.get_dag(dag_id)
         assert dag is not None
         assert dag_id == dag.dag_id
-        assert dagbag.process_file_calls == 2
+        assert dagbag.import_calls == 2
 
     @patch.object(DagModel, "get_current")
     def test_refresh_packaged_dag(self, mock_dagmodel, test_zip_path):
@@ -950,20 +1050,20 @@ class TestDagBag:
         mock_dagmodel.return_value.fileloc = fileloc
 
         class _TestDagBag(DagBag):
-            process_file_calls = 0
+            import_calls = 0
 
-            def process_file(self, filepath, only_if_updated=True, 
safe_mode=True):
-                if filepath in fileloc:
-                    _TestDagBag.process_file_calls += 1
-                return super().process_file(filepath, only_if_updated, 
safe_mode)
+            def _process_definition(self, importer, definition, *, 
only_if_updated):
+                if repr(definition) == fileloc:
+                    _TestDagBag.import_calls += 1
+                return super()._process_definition(importer, definition, 
only_if_updated=only_if_updated)
 
         dagbag = _TestDagBag(dag_folder=os.path.realpath(test_zip_path))
 
-        assert dagbag.process_file_calls == 1
+        assert dagbag.import_calls == 1
         dag = dagbag.get_dag(dag_id)
         assert dag is not None
         assert dag_id == dag.dag_id
-        assert dagbag.process_file_calls == 2
+        assert dagbag.import_calls == 2
 
     def process_dag(self, create_dag, tmp_path):
         """
@@ -1199,6 +1299,8 @@ with airflow.DAG(
                 f"{dag_file}:48: UserWarning: Some Warning",
             )
         }
+        # Python warning categories are not Dag warnings
+        assert dagbag.dag_warnings == set()
 
         with warnings.catch_warnings():
             # Disable capture DeprecationWarning, and it should be reflected 
in captured warnings
@@ -1227,7 +1329,7 @@ with airflow.DAG(
         dagbag = DagBag(dag_folder=warning_zipped_dag_path)
         assert dagbag.dagbag_stats[0].warning_num == 2
         assert dagbag.captured_warnings == {
-            warning_zipped_dag_path: (
+            in_zip_dag_file: (
                 f"{in_zip_dag_file}:46: DeprecationWarning: Deprecated 
Parameter",
                 f"{in_zip_dag_file}:48: UserWarning: Some Warning",
             )
@@ -1264,6 +1366,25 @@ with airflow.DAG(
         dagbag.bag_dag(dag)
         assert dagbag.dag_warnings == expected
 
+    @conf_vars({("dag_processor", "dag_importer_configs"): 
json.dumps([f"{__name__}.TextDagImporter"])})
+    def test_custom_importer_dags_are_bound_to_their_definition(self, 
tmp_path):
+        dag_file = tmp_path / "nested" / "text_dag.dagtxt"
+        dag_file.parent.mkdir()
+        dag_file.write_text("text_dag\n")
+
+        dagbag = DagBag(dag_folder=tmp_path, bundle_path=tmp_path, 
bundle_name="test_bundle")
+
+        dag = dagbag.dags["text_dag"]
+        assert dag.fileloc == os.fspath(dag_file)
+        assert dag.relative_fileloc == "nested/text_dag.dagtxt"
+        assert dagbag.dag_source_codes == {
+            os.fspath(dag_file): DagSourceCode(source_code="text_dag\n", 
language="text")
+        }
+        assert dagbag.dag_warnings == {DagWarning("text_dag", 
"test:deprecated_field", "Deprecated field")}
+        assert dagbag.captured_warnings == {
+            os.fspath(dag_file): (f"{dag_file}:1: test:deprecated_field: 
Deprecated field",)
+        }
+
     def test_sigsegv_handling(self, tmp_path, caplog):
         """
         Test that a SIGSEGV in a DAG file is handled gracefully and does not 
crash the process.
@@ -1312,68 +1433,12 @@ with airflow.DAG(
                 """
             )
         )
-        with 
mock.patch("airflow.dag_processing.importers.python_importer.signal.signal") as 
mock_signal:
+        with mock.patch("airflow.sdk.importers.python_importer.signal.signal") 
as mock_signal:
             mock_signal.side_effect = ValueError("Invalid signal setting")
             DagBag(dag_folder=os.fspath(tmp_path))
             assert "SIGSEGV signal handler registration failed. Not in the 
main thread" in caplog.text
 
 
-class TestCaptureWithReraise:
-    @staticmethod
-    def raise_warnings():
-        warnings.warn("Foo", UserWarning, stacklevel=2)
-        warnings.warn("Bar", UserWarning, stacklevel=2)
-        warnings.warn("Baz", UserWarning, stacklevel=2)
-
-    def test_capture_no_warnings(self):
-        with warnings.catch_warnings():
-            warnings.simplefilter("error")
-            with _capture_with_reraise() as cw:
-                pass
-            assert cw == []
-
-    def test_capture_warnings(self):
-        with pytest.warns(UserWarning, match="(Foo|Bar|Baz)") as ctx:
-            with _capture_with_reraise() as cw:
-                self.raise_warnings()
-        assert len(cw) == 3
-        assert len(ctx.list) == 3
-
-    def test_capture_warnings_with_parent_error_filter(self):
-        with warnings.catch_warnings(record=True) as records:
-            warnings.filterwarnings("error", message="Bar")
-            with _capture_with_reraise() as cw:
-                with pytest.raises(UserWarning, match="Bar"):
-                    self.raise_warnings()
-            assert len(cw) == 1
-        assert len(records) == 1
-
-    def test_capture_warnings_with_parent_ignore_filter(self):
-        with warnings.catch_warnings(record=True) as records:
-            warnings.filterwarnings("ignore", message="Baz")
-            with _capture_with_reraise() as cw:
-                self.raise_warnings()
-            assert len(cw) == 2
-        assert len(records) == 2
-
-    def test_capture_warnings_with_filters(self):
-        with warnings.catch_warnings(record=True) as records:
-            with _capture_with_reraise() as cw:
-                warnings.filterwarnings("ignore", message="Foo")
-                self.raise_warnings()
-            assert len(cw) == 2
-        assert len(records) == 2
-
-    def test_capture_warnings_with_error_filters(self):
-        with warnings.catch_warnings(record=True) as records:
-            with _capture_with_reraise() as cw:
-                warnings.filterwarnings("error", message="Bar")
-                with pytest.raises(UserWarning, match="Bar"):
-                    self.raise_warnings()
-            assert len(cw) == 1
-        assert len(records) == 1
-
-
 class TestBundlePathSysPath:
     """Tests for bundle_path sys.path handling in BundleDagBag."""
 
diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py 
b/airflow-core/tests/unit/dag_processing/test_manager.py
index b581a3d742d..a3624e1d4a1 100644
--- a/airflow-core/tests/unit/dag_processing/test_manager.py
+++ b/airflow-core/tests/unit/dag_processing/test_manager.py
@@ -72,6 +72,7 @@ from airflow.models.serialized_dag import SerializedDagModel
 from airflow.models.team import Team
 from airflow.providers.standard.operators.empty import EmptyOperator
 from airflow.sdk import DAG as SdkDAG
+from airflow.sdk.importers import DagSourceCode
 from airflow.serialization.serialized_objects import LazyDeserializedDAG
 from airflow.utils.net import get_hostname
 from airflow.utils.session import create_session
@@ -2090,6 +2091,30 @@ class TestDagFileProcessorManager:
         # and the DAG from test_dag2.py is deactivated
         assert session.get(DagModel, "test_dag2").is_stale is True
 
+    
@mock.patch("airflow.dag_processing.manager.update_dag_parsing_results_in_db", 
autospec=True)
+    def 
test_persist_parsing_result_passes_parsed_definitions_and_source_codes(self, 
mock_update):
+        source_codes = {"/bundle/dags.zip/a.py": 
DagSourceCode(source_code="src", language="python")}
+        parsing_result = DagFileParsingResult(
+            fileloc="/bundle/dags.zip",
+            serialized_dags=[],
+            parsed_definitions=["dags.zip/a.py"],
+            dag_source_codes=source_codes,
+        )
+
+        DagFileProcessorManager(max_runs=1).persist_parsing_result(
+            bundle_name="testing",
+            bundle_version=None,
+            version_data=None,
+            parsing_result=parsing_result,
+            run_duration=1.0,
+            relative_fileloc="dags.zip",
+            session=mock.sentinel.session,
+        )
+
+        kwargs = mock_update.call_args.kwargs
+        assert kwargs["files_parsed"] == {("testing", "dags.zip"), ("testing", 
"dags.zip/a.py")}
+        assert kwargs["dag_source_codes"] == source_codes
+
     @pytest.mark.parametrize(
         ("rel_filelocs", "expected_return", "expected_dag1_stale", 
"expected_dag2_stale"),
         [
diff --git a/airflow-core/tests/unit/dag_processing/test_processor.py 
b/airflow-core/tests/unit/dag_processing/test_processor.py
index 57cce320a80..0cb61aa0286 100644
--- a/airflow-core/tests/unit/dag_processing/test_processor.py
+++ b/airflow-core/tests/unit/dag_processing/test_processor.py
@@ -19,15 +19,17 @@ from __future__ import annotations
 
 import inspect
 import logging
+import os
 import pathlib
 import sys
 import textwrap
 import typing
 import uuid
+import zipfile
 from collections.abc import Callable, Iterable
 from socket import socketpair
 from typing import TYPE_CHECKING, Any, BinaryIO
-from unittest.mock import MagicMock, patch
+from unittest.mock import MagicMock, PropertyMock, patch
 
 import pytest
 import structlog
@@ -66,6 +68,7 @@ from airflow.dag_processing.processor import (
     _pre_import_airflow_modules,
 )
 from airflow.models import DagRun
+from airflow.models.dagwarning import DagWarning
 from airflow.sdk import DAG, BaseOperator
 from airflow.sdk.api.client import Client
 from airflow.sdk.api.datamodels._generated import ConnectionResponse, 
DagRunState, VariableResponse
@@ -85,6 +88,7 @@ from airflow.sdk.execution_time.comms import (
     XComSequenceSliceResult,
 )
 from airflow.sdk.execution_time.task_runner import RuntimeTaskInstance
+from airflow.sdk.importers import DagSourceCode
 from airflow.utils.session import create_session
 from airflow.utils.state import TaskInstanceState
 
@@ -767,6 +771,29 @@ def test_parse_file_static_check_with_default_warning():
     )
 
 
[email protected](DagBag, "dag_warnings", new_callable=PropertyMock)
+def 
test_parse_file_reports_definitions_source_codes_and_dag_warnings(mock_dag_warnings,
 tmp_path):
+    source = "from airflow.sdk import DAG\n\ndag = DAG('member', 
schedule=None)\n"
+    with zipfile.ZipFile(tmp_path / "dags.zip", "w") as zf:
+        zf.writestr("member.py", source)
+    mock_dag_warnings.return_value = {DagWarning("member", 
"test:deprecated_field", "Deprecated field")}
+
+    result = _parse_file(
+        DagFileParseRequest(
+            file=os.fspath(tmp_path / "dags.zip"), bundle_path=tmp_path, 
bundle_name="testing"
+        ),
+        log=structlog.get_logger(),
+    )
+
+    assert result.parsed_definitions == ["dags.zip/member.py"]
+    assert result.dag_source_codes == {
+        os.fspath(tmp_path / "dags.zip" / "member.py"): 
DagSourceCode(source_code=source, language="python")
+    }
+    assert result.warnings == [
+        {"dag_id": "member", "warning_type": "test:deprecated_field", 
"message": "Deprecated field"}
+    ]
+
+
 def test_callback_processing_does_not_update_timestamps():
     """Callback processing should not update last_finish_time to prevent stale 
DAG detection."""
     stat = process_parse_results(
@@ -2139,6 +2166,8 @@ class TestExecuteEmailCallbacks:
             mock_dagbag_instance = MagicMock()
             mock_dagbag_instance.dags = {}
             mock_dagbag_instance.import_errors = {}  # Must be a dict, not 
MagicMock for Pydantic validation
+            mock_dagbag_instance.parsed_definitions = []
+            mock_dagbag_instance.dag_source_codes = {}
             mock_dagbag_class.return_value = mock_dagbag_instance
 
             request = DagFileParseRequest(
diff --git a/airflow-core/tests/unit/models/test_dagcode.py 
b/airflow-core/tests/unit/models/test_dagcode.py
index 419a8f0cbd5..c5276d85f10 100644
--- a/airflow-core/tests/unit/models/test_dagcode.py
+++ b/airflow-core/tests/unit/models/test_dagcode.py
@@ -29,6 +29,7 @@ from airflow.dag_processing.dagbag import DagBag
 from airflow.models.dag_version import DagVersion
 from airflow.models.dagcode import DagCode
 from airflow.sdk import task as task_decorator
+from airflow.sdk.importers import DagSourceCode
 from airflow.serialization.definitions.dag import SerializedDAG
 
 # To move it to a shared module.
@@ -270,3 +271,41 @@ class TestDagCode:
         sync_dag_to_db(dag)
 
         assert DagCode.get_latest_dagcode(dag.dag_id).language == "python"
+
+    def test_write_code_with_dag_source_code(self, dag_maker, session):
+        with dag_maker("dag_source_code_test") as dag:
+            pass
+        sync_dag_to_db(dag)
+        dag_version = DagVersion.get_latest_version(dag.dag_id)
+        clear_db_dag_code()
+
+        dag_code = DagCode.write_code(
+            dag_version,
+            dag.fileloc,
+            dag_source_code=DagSourceCode(source_code="dag: my_dag\nversion: 
1", language="yaml"),
+            session=session,
+        )
+        session.commit()
+
+        stored = DagCode.get_latest_dagcode(dag.dag_id, session=session)
+        assert (stored.id, stored.source_code, stored.language) == (
+            dag_code.id,
+            "dag: my_dag\nversion: 1",
+            "yaml",
+        )
+
+    def test_update_source_code_with_dag_source_code(self, dag_maker, session):
+        with dag_maker("dag_source_code_update") as dag:
+            pass
+        sync_dag_to_db(dag)
+
+        DagCode.update_source_code(
+            dag.dag_id,
+            dag.fileloc,
+            dag_source_code=DagSourceCode(source_code="print('custom 
importer')", language="custom_lang"),
+            session=session,
+        )
+        session.commit()
+
+        latest = DagCode.get_latest_dagcode(dag.dag_id, session=session)
+        assert (latest.source_code, latest.language) == ("print('custom 
importer')", "custom_lang")
diff --git a/devel-common/src/tests_common/test_utils/config.py 
b/devel-common/src/tests_common/test_utils/config.py
index 619a88a6930..e5724ea6525 100644
--- a/devel-common/src/tests_common/test_utils/config.py
+++ b/devel-common/src/tests_common/test_utils/config.py
@@ -19,6 +19,7 @@ from __future__ import annotations
 
 import contextlib
 import os
+import sys
 from typing import TYPE_CHECKING, Literal, overload
 
 if TYPE_CHECKING:
@@ -66,8 +67,6 @@ PROVIDER_METADATA_OVERRIDES_CFG_FALLBACK: list[tuple[str, 
str, str, str]] = [
 @contextlib.contextmanager
 def conf_vars(overrides):
     """Automatically detects which config modules are loaded (Core, SDK, or 
both) and updates them accordingly temporarily."""
-    import sys
-
     from airflow import settings
 
     configs = []
@@ -99,6 +98,7 @@ def conf_vars(overrides):
     if "airflow.configuration" in sys.modules:
         settings.configure_vars()
     _clear_dag_bundle_config_cache()
+    _clear_importer_registry_cache()
 
     try:
         yield
@@ -118,6 +118,7 @@ def conf_vars(overrides):
         if "airflow.configuration" in sys.modules:
             settings.configure_vars()
         _clear_dag_bundle_config_cache()
+        _clear_importer_registry_cache()
 
 
 def _clear_dag_bundle_config_cache() -> None:
@@ -131,6 +132,12 @@ def _clear_dag_bundle_config_cache() -> None:
         cache.cache_clear()
 
 
+def _clear_importer_registry_cache() -> None:
+    """Drop the cached Dag importer registries (and their importers) so they 
read the current config."""
+    if importers := sys.modules.get("airflow.sdk.importers.base"):
+        importers.reset_importer_registry()
+
+
 @overload
 def create_fresh_airflow_config(variant: Literal["core"] = ...) -> 
AirflowConfigParser: ...
 
diff --git a/generated/known_sdk_imports_in_core.txt 
b/generated/known_sdk_imports_in_core.txt
index f05d5091514..b906ba25855 100644
--- a/generated/known_sdk_imports_in_core.txt
+++ b/generated/known_sdk_imports_in_core.txt
@@ -3,11 +3,9 @@ 
airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_04_06.py::1
 airflow-core/src/airflow/cli/commands/task_command.py::7
 airflow-core/src/airflow/cli/commands/triggerer_command.py::1
 airflow-core/src/airflow/configuration.py::1
-airflow-core/src/airflow/dag_processing/dagbag.py::1
-airflow-core/src/airflow/dag_processing/importers/base.py::1
-airflow-core/src/airflow/dag_processing/importers/python_importer.py::7
+airflow-core/src/airflow/dag_processing/dagbag.py::3
 airflow-core/src/airflow/dag_processing/manager.py::4
-airflow-core/src/airflow/dag_processing/processor.py::13
+airflow-core/src/airflow/dag_processing/processor.py::14
 airflow-core/src/airflow/exceptions.py::1
 airflow-core/src/airflow/executors/base_executor.py::3
 airflow-core/src/airflow/jobs/triggerer_job_runner.py::18
diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json 
b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json
index 4777a686c98..d6ee95927ef 100644
--- a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json
+++ b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json
@@ -1015,6 +1015,20 @@
           "default": null,
           "title": "Import Errors"
         },
+        "parsed_definitions": {
+          "items": {
+            "type": "string"
+          },
+          "title": "Parsed Definitions",
+          "type": "array"
+        },
+        "dag_source_codes": {
+          "additionalProperties": {
+            "$ref": "#/$defs/DagSourceCode"
+          },
+          "title": "Dag Source Codes",
+          "type": "object"
+        },
         "type": {
           "const": "DagFileParsingResult",
           "default": "DagFileParsingResult",
@@ -1493,6 +1507,25 @@
       "title": "DagRunType",
       "type": "string"
     },
+    "DagSourceCode": {
+      "description": "Raw source code and its language identifier for a DAG 
definition.",
+      "properties": {
+        "source_code": {
+          "title": "Source Code",
+          "type": "string"
+        },
+        "language": {
+          "title": "Language",
+          "type": "string"
+        }
+      },
+      "required": [
+        "source_code",
+        "language"
+      ],
+      "title": "DagSourceCode",
+      "type": "object"
+    },
     "DeferTask": {
       "additionalProperties": false,
       "description": "Update a task instance state to deferred.",
diff --git 
a/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py 
b/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py
index 06be64d3461..c974f797bb0 100644
--- a/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py
+++ b/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py
@@ -39,12 +39,18 @@ def get_bundle() -> VersionBundle:
 
     from airflow.sdk.execution_time.schema.versions.v2026_10_30 import (
         AddArgBindingsToSupervisorTIRunContext,
+        AddDagDefinitionsToDagFileParsingResult,
         AddRetryReasonToTaskState,
     )
 
     return VersionBundle(
         HeadVersion(),
-        Version("2026-10-30", AddArgBindingsToSupervisorTIRunContext, 
AddRetryReasonToTaskState),
+        Version(
+            "2026-10-30",
+            AddArgBindingsToSupervisorTIRunContext,
+            AddRetryReasonToTaskState,
+            AddDagDefinitionsToDagFileParsingResult,
+        ),
         Version("2026-06-16"),
     )
 
diff --git 
a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py 
b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py
index af54445542b..4aca18cf50c 100644
--- a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py
+++ b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_10_30.py
@@ -28,6 +28,7 @@ from __future__ import annotations
 
 from cadwyn import VersionChange, schema
 
+from airflow.dag_processing.processor import DagFileParsingResult  # noqa: 
SDK002
 from airflow.sdk.api.datamodels._generated import TIRunContext
 from airflow.sdk.execution_time.comms import TaskState
 
@@ -52,3 +53,14 @@ class AddRetryReasonToTaskState(VersionChange):
     description = __doc__
 
     instructions_to_migrate_to_previous_version = 
(schema(TaskState).field("retry_reason").didnt_exist,)
+
+
+class AddDagDefinitionsToDagFileParsingResult(VersionChange):
+    """Add the imported Dag definitions and their source code to 
`DagFileParsingResult`."""
+
+    description = __doc__
+
+    instructions_to_migrate_to_previous_version = (
+        schema(DagFileParsingResult).field("parsed_definitions").didnt_exist,
+        schema(DagFileParsingResult).field("dag_source_codes").didnt_exist,
+    )
diff --git a/task-sdk/src/airflow/sdk/importers/base.py 
b/task-sdk/src/airflow/sdk/importers/base.py
index 94916fa60eb..82798727fb2 100644
--- a/task-sdk/src/airflow/sdk/importers/base.py
+++ b/task-sdk/src/airflow/sdk/importers/base.py
@@ -229,6 +229,8 @@ class AbstractDagImporter(ABC, Generic[DefT]):
     :meth:`.list_dag_definitions` yields definitions of :class:`DagDefinition`
     subtypes, and those same objects are fed back to :meth:`import_definition`,
     so a concrete importer only ever deals with its own definition type.
+
+    .. note:: |experimental|
     """
 
     @abstractmethod
diff --git a/ts-sdk/src/generated/supervisor.ts 
b/ts-sdk/src/generated/supervisor.ts
index 1f9fe56f892..76e9e727815 100644
--- a/ts-sdk/src/generated/supervisor.ts
+++ b/ts-sdk/src/generated/supervisor.ts
@@ -277,6 +277,9 @@ export type Warnings = unknown[] | null;
 export type ImportErrors = {
   [k: string]: string;
 } | null;
+export type ParsedDefinitions = string[];
+export type SourceCode = string;
+export type Language = string;
 export type Type16 = "DagFileParsingResult";
 export type DagId4 = string;
 export type IsPaused = boolean;
@@ -1112,6 +1115,8 @@ export interface DagFileParsingResult {
   serialized_dags: SerializedDags;
   warnings?: Warnings;
   import_errors?: ImportErrors;
+  parsed_definitions?: ParsedDefinitions;
+  dag_source_codes?: DagSourceCodes;
   type?: Type16;
 }
 /**
@@ -1130,6 +1135,19 @@ export interface LazyDeserializedDAG {
 export interface Data {
   [k: string]: unknown;
 }
+export interface DagSourceCodes {
+  [k: string]: DagSourceCode;
+}
+/**
+ * Raw source code and its language identifier for a DAG definition.
+ *
+ * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
+ * via the `definition` "DagSourceCode".
+ */
+export interface DagSourceCode {
+  source_code: SourceCode;
+  language: Language;
+}
 /**
  * This interface was referenced by `SupervisorWireSchema`'s JSON-Schema
  * via the `definition` "DagResult".

Reply via email to