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