This is an automated email from the ASF dual-hosted git repository.
shahar1 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 30d0d7bac4a Adjust warnings, import errors, and Dag source models for
sdk/importers (#72474)
30d0d7bac4a is described below
commit 30d0d7bac4a74f6a1eb65947d070bf1418ca4cc9
Author: Dilnaz Amanzholova <[email protected]>
AuthorDate: Mon Sep 28 12:07:46 2026 +0200
Adjust warnings, import errors, and Dag source models for sdk/importers
(#72474)
---
airflow-core/docs/migrations-ref.rst | 6 +-
airflow-core/src/airflow/api/common/delete_dag.py | 7 +-
.../api_fastapi/core_api/datamodels/dag_sources.py | 1 +
.../api_fastapi/core_api/datamodels/dag_warning.py | 4 +-
.../core_api/datamodels/import_error.py | 1 +
.../core_api/openapi/v2-rest-api-generated.yaml | 25 ++++++-
.../core_api/routes/public/dag_sources.py | 1 +
.../core_api/routes/public/dag_warning.py | 6 +-
.../core_api/routes/public/import_error.py | 2 +
.../src/airflow/dag_processing/collection.py | 2 +
..._3_4_0_add_source_reference_to_import_error.py} | 41 ++++++-----
.../0140_3_4_0_add_language_to_dag_code.py} | 43 ++++++-----
airflow-core/src/airflow/models/dagcode.py | 3 +
airflow-core/src/airflow/models/dagwarning.py | 31 +++++++-
airflow-core/src/airflow/models/errors.py | 12 ++-
.../src/airflow/ui/openapi-gen/queries/common.ts | 4 +-
.../ui/openapi-gen/queries/ensureQueryData.ts | 6 +-
.../src/airflow/ui/openapi-gen/queries/prefetch.ts | 6 +-
.../src/airflow/ui/openapi-gen/queries/queries.ts | 6 +-
.../src/airflow/ui/openapi-gen/queries/suspense.ts | 6 +-
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 36 ++++++++-
.../ui/openapi-gen/requests/services.gen.ts | 2 +-
.../airflow/ui/openapi-gen/requests/types.gen.ts | 8 +-
.../src/airflow/ui/src/pages/Dag/Code/Code.tsx | 10 ++-
airflow-core/src/airflow/utils/db.py | 2 +-
.../tests/unit/api/common/test_delete_dag.py | 49 ++++++++++++
.../core_api/routes/public/test_dag_sources.py | 4 +
.../core_api/routes/public/test_dag_warning.py | 19 +++++
.../core_api/routes/public/test_import_error.py | 76 ++++++++++++++++++-
airflow-core/tests/unit/models/test_dagcode.py | 8 ++
airflow-core/tests/unit/models/test_dagwarning.py | 38 +++++++++-
airflow-core/tests/unit/models/test_errors.py | 86 ++++++++++++++++++++++
.../src/airflowctl/api/datamodels/generated.py | 10 ++-
.../tests/airflow_ctl/api/test_operations.py | 1 +
34 files changed, 483 insertions(+), 79 deletions(-)
diff --git a/airflow-core/docs/migrations-ref.rst
b/airflow-core/docs/migrations-ref.rst
index c0b456561fd..4f5bf0f33ad 100644
--- a/airflow-core/docs/migrations-ref.rst
+++ b/airflow-core/docs/migrations-ref.rst
@@ -39,7 +39,11 @@ Here's the list of all the Database Migrations that are
executed via when you ru
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| Revision ID | Revises ID | Airflow Version | Description
|
+=========================+==================+===================+==============================================================+
-| ``a61f0c9d2b47`` (head) | ``c9f4b3e7a218`` | ``3.4.0`` | Allocate
attempt numbers for pending retries and cleared |
+| ``e5a91c7f42b3`` (head) | ``ca8499dc1004`` | ``3.4.0`` | Add
language column to dag_code. |
++-------------------------+------------------+-------------------+--------------------------------------------------------------+
+| ``ca8499dc1004`` | ``a61f0c9d2b47`` | ``3.4.0`` | Add
source_reference to import_error. |
++-------------------------+------------------+-------------------+--------------------------------------------------------------+
+| ``a61f0c9d2b47`` | ``c9f4b3e7a218`` | ``3.4.0`` | Allocate
attempt numbers for pending retries and cleared |
| | | | tasks.
|
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``c9f4b3e7a218`` | ``3b7a91c5df20`` | ``3.4.0`` | Replace the
``job.team_name`` column with a ``job_team`` |
diff --git a/airflow-core/src/airflow/api/common/delete_dag.py
b/airflow-core/src/airflow/api/common/delete_dag.py
index e32e8bb90f3..4dc753fa945 100644
--- a/airflow-core/src/airflow/api/common/delete_dag.py
+++ b/airflow-core/src/airflow/api/common/delete_dag.py
@@ -22,7 +22,7 @@ from __future__ import annotations
import logging
from typing import TYPE_CHECKING, cast
-from sqlalchemy import delete, select
+from sqlalchemy import delete, or_, select
from airflow import models
from airflow.exceptions import AirflowException, DagNotFound
@@ -82,7 +82,10 @@ def delete_dag(dag_id: str, keep_records_in_log: bool =
True, *, session: Sessio
# This handles the case when the dag_id is changed in the file
session.execute(
delete(ParseImportError).where(
- ParseImportError.filename == dag.relative_fileloc,
+ or_(
+ ParseImportError.source_reference == dag.relative_fileloc,
+ ParseImportError.filename == dag.relative_fileloc,
+ ),
ParseImportError.bundle_name == dag.bundle_name,
)
)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_sources.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_sources.py
index fdc081a0ab2..d6dc5a3fc91 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_sources.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_sources.py
@@ -28,3 +28,4 @@ class DAGSourceResponse(BaseModel):
dag_id: str
version_number: int | None
dag_display_name: str = Field(validation_alias=AliasPath("dag_model",
"dag_display_name"))
+ language: str | None = None
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
index 470a0b14fc6..24dd095c75a 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
@@ -23,14 +23,14 @@ from datetime import datetime
from pydantic import AliasPath, Field
from airflow.api_fastapi.core_api.base import BaseModel
-from airflow.models.dagwarning import DagWarningType
+from airflow.models.dagwarning import DagWarningTypeValue
class DAGWarningResponse(BaseModel):
"""Dag Warning serializer for responses."""
dag_id: str
- warning_type: DagWarningType
+ warning_type: DagWarningTypeValue
message: str
timestamp: datetime
dag_display_name: str = Field(validation_alias=AliasPath("dag_model",
"dag_display_name"))
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
index d070df98fdb..17c4e2e127b 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/import_error.py
@@ -31,6 +31,7 @@ class ImportErrorResponse(BaseModel):
id: int = Field(alias="import_error_id")
timestamp: datetime
filename: str
+ source_reference: str | None
bundle_name: str | None
stacktrace: str = Field(alias="stack_trace")
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index d57c05b6812..dabcc47c9b7 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -3805,6 +3805,9 @@ paths:
schema:
anyOf:
- $ref: '#/components/schemas/DagWarningType'
+ - type: string
+ maxLength: 50
+ pattern: ^[a-z][a-z0-9_]*:[a-z0-9_.\-]+$
- type: 'null'
title: Warning Type
- name: limit
@@ -5256,13 +5259,13 @@ paths:
type: string
description: 'Attributes to order by, multi criteria sort is
supported.
Prefix with `-` for descending order. Supported attributes: `id,
timestamp,
- filename, bundle_name, stacktrace, import_error_id`'
+ filename, source_reference, bundle_name, stacktrace,
import_error_id`'
default:
- id
title: Order By
description: 'Attributes to order by, multi criteria sort is
supported. Prefix
with `-` for descending order. Supported attributes: `id, timestamp,
filename,
- bundle_name, stacktrace, import_error_id`'
+ source_reference, bundle_name, stacktrace, import_error_id`'
- name: filename_pattern
in: query
required: false
@@ -14283,6 +14286,11 @@ components:
dag_display_name:
type: string
title: Dag Display Name
+ language:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Language
type: object
required:
- content
@@ -14345,7 +14353,12 @@ components:
type: string
title: Dag Id
warning_type:
- $ref: '#/components/schemas/DagWarningType'
+ anyOf:
+ - $ref: '#/components/schemas/DagWarningType'
+ - type: string
+ maxLength: 50
+ pattern: ^[a-z][a-z0-9_]*:[a-z0-9_.\-]+$
+ title: Warning Type
message:
type: string
title: Message
@@ -15450,6 +15463,11 @@ components:
filename:
type: string
title: Filename
+ source_reference:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Source Reference
bundle_name:
anyOf:
- type: string
@@ -15469,6 +15487,7 @@ components:
- import_error_id
- timestamp
- filename
+ - source_reference
- bundle_name
- stack_trace
- file_token
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
index f74ce76c7b7..212e63c55f2 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
@@ -103,6 +103,7 @@ def get_dag_source(
content=content,
version_number=dag_version.version_number,
dag_display_name=dag_version.dag_model.dag_display_name,
+ language=dag_version.dag_code.language,
)
if accept == Mimetype.TEXT:
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_warning.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_warning.py
index c1a652c370f..0096b82d1ba 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_warning.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_warning.py
@@ -40,7 +40,7 @@ from airflow.api_fastapi.core_api.datamodels.dag_warning
import (
DAGWarningCollectionResponse,
)
from airflow.api_fastapi.core_api.security import
ReadableDagWarningsFilterDep, requires_access_dag
-from airflow.models.dagwarning import DagWarning, DagWarningType
+from airflow.models.dagwarning import DagWarning, DagWarningTypeValue
dag_warning_router = AirflowRouter(tags=["DagWarning"])
@@ -52,8 +52,8 @@ dag_warning_router = AirflowRouter(tags=["DagWarning"])
def list_dag_warnings(
dag_id: Annotated[FilterParam[str | None],
Depends(filter_param_factory(DagWarning.dag_id, str | None))],
warning_type: Annotated[
- FilterParam[DagWarningType | None],
- Depends(filter_param_factory(DagWarning.warning_type, DagWarningType |
None)),
+ FilterParam[DagWarningTypeValue | None],
+ Depends(filter_param_factory(DagWarning.warning_type,
DagWarningTypeValue | None)),
],
limit: QueryLimit,
offset: QueryOffset,
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/import_error.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/import_error.py
index d5801e6b7e4..fe8cb80a692 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/import_error.py
+++
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/import_error.py
@@ -147,6 +147,7 @@ def get_import_errors(
"id",
"timestamp",
"filename",
+ "source_reference",
"bundle_name",
"stacktrace",
],
@@ -263,6 +264,7 @@ def get_import_errors(
ParseImportError.id,
ParseImportError.timestamp,
ParseImportError.filename,
+ ParseImportError.source_reference,
ParseImportError.bundle_name,
ParseImportError.stacktrace,
)
diff --git a/airflow-core/src/airflow/dag_processing/collection.py
b/airflow-core/src/airflow/dag_processing/collection.py
index 36b948abad6..70c95f8eedf 100644
--- a/airflow-core/src/airflow/dag_processing/collection.py
+++ b/airflow-core/src/airflow/dag_processing/collection.py
@@ -424,6 +424,7 @@ def _update_import_errors(
)
.values(
filename=relative_fileloc,
+ source_reference=relative_fileloc,
bundle_name=bundle_name_,
timestamp=utcnow(),
stacktrace=stacktrace,
@@ -447,6 +448,7 @@ def _update_import_errors(
else:
import_error = ParseImportError(
filename=relative_fileloc,
+ source_reference=relative_fileloc,
bundle_name=bundle_name_,
timestamp=utcnow(),
stacktrace=stacktrace,
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
b/airflow-core/src/airflow/migrations/versions/0139_3_4_0_add_source_reference_to_import_error.py
similarity index 51%
copy from
airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
copy to
airflow-core/src/airflow/migrations/versions/0139_3_4_0_add_source_reference_to_import_error.py
index 470a0b14fc6..42613b825d3 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
+++
b/airflow-core/src/airflow/migrations/versions/0139_3_4_0_add_source_reference_to_import_error.py
@@ -1,3 +1,4 @@
+#
# 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
@@ -15,29 +16,35 @@
# specific language governing permissions and limitations
# under the License.
-from __future__ import annotations
+"""
+Add source_reference to import_error.
-from collections.abc import Iterable
-from datetime import datetime
+Revision ID: ca8499dc1004
+Revises: a61f0c9d2b47
+Create Date: 2026-08-27 12:45:02.276898
-from pydantic import AliasPath, Field
+"""
-from airflow.api_fastapi.core_api.base import BaseModel
-from airflow.models.dagwarning import DagWarningType
+from __future__ import annotations
+import sqlalchemy as sa
+from alembic import op
-class DAGWarningResponse(BaseModel):
- """Dag Warning serializer for responses."""
+# revision identifiers, used by Alembic.
+revision = "ca8499dc1004"
+down_revision = "a61f0c9d2b47"
+branch_labels = None
+depends_on = None
+airflow_version = "3.4.0"
- dag_id: str
- warning_type: DagWarningType
- message: str
- timestamp: datetime
- dag_display_name: str = Field(validation_alias=AliasPath("dag_model",
"dag_display_name"))
+def upgrade():
+ """Apply add source_reference to import_error."""
+ with op.batch_alter_table("import_error", schema=None) as batch_op:
+ batch_op.add_column(sa.Column("source_reference",
sa.String(length=2000), nullable=True))
-class DAGWarningCollectionResponse(BaseModel):
- """Dag warning collection serializer for responses."""
- dag_warnings: Iterable[DAGWarningResponse]
- total_entries: int
+def downgrade():
+ """Unapply add source_reference to import_error."""
+ with op.batch_alter_table("import_error", schema=None) as batch_op:
+ batch_op.drop_column("source_reference")
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
b/airflow-core/src/airflow/migrations/versions/0140_3_4_0_add_language_to_dag_code.py
similarity index 51%
copy from
airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
copy to
airflow-core/src/airflow/migrations/versions/0140_3_4_0_add_language_to_dag_code.py
index 470a0b14fc6..9672fdec2a6 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_warning.py
+++
b/airflow-core/src/airflow/migrations/versions/0140_3_4_0_add_language_to_dag_code.py
@@ -1,3 +1,4 @@
+#
# 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
@@ -15,29 +16,37 @@
# specific language governing permissions and limitations
# under the License.
-from __future__ import annotations
+"""
+Add language column to dag_code.
-from collections.abc import Iterable
-from datetime import datetime
+Revision ID: e5a91c7f42b3
+Revises: ca8499dc1004
+Create Date: 2026-09-10 12:00:00.000000
-from pydantic import AliasPath, Field
+"""
-from airflow.api_fastapi.core_api.base import BaseModel
-from airflow.models.dagwarning import DagWarningType
+from __future__ import annotations
+import sqlalchemy as sa
+from alembic import op
-class DAGWarningResponse(BaseModel):
- """Dag Warning serializer for responses."""
+# revision identifiers, used by Alembic.
+revision = "e5a91c7f42b3"
+down_revision = "ca8499dc1004"
+branch_labels = None
+depends_on = None
+airflow_version = "3.4.0"
- dag_id: str
- warning_type: DagWarningType
- message: str
- timestamp: datetime
- dag_display_name: str = Field(validation_alias=AliasPath("dag_model",
"dag_display_name"))
+def upgrade():
+ """Apply add language column to dag_code."""
+ with op.batch_alter_table("dag_code", schema=None) as batch_op:
+ batch_op.add_column(
+ sa.Column("language", sa.String(length=64), nullable=False,
server_default="python")
+ )
-class DAGWarningCollectionResponse(BaseModel):
- """Dag warning collection serializer for responses."""
- dag_warnings: Iterable[DAGWarningResponse]
- total_entries: int
+def downgrade():
+ """Unapply add language column to dag_code."""
+ with op.batch_alter_table("dag_code", schema=None) as batch_op:
+ batch_op.drop_column("language")
diff --git a/airflow-core/src/airflow/models/dagcode.py
b/airflow-core/src/airflow/models/dagcode.py
index b0306b15ff1..b21af49d57f 100644
--- a/airflow-core/src/airflow/models/dagcode.py
+++ b/airflow-core/src/airflow/models/dagcode.py
@@ -66,6 +66,9 @@ class DagCode(Base):
)
source_code: Mapped[str] = mapped_column(Text().with_variant(MEDIUMTEXT(),
"mysql"), nullable=False)
source_code_hash: Mapped[str] = mapped_column(String(32), nullable=False)
+ language: Mapped[str] = mapped_column(
+ String(64), nullable=False, default="python", server_default="python"
+ )
dag_version_id: Mapped[UUID] = mapped_column(
sa.Uuid(), ForeignKey("dag_version.id", ondelete="CASCADE"),
nullable=False, unique=True
)
diff --git a/airflow-core/src/airflow/models/dagwarning.py
b/airflow-core/src/airflow/models/dagwarning.py
index b0a5a3c93e1..6305361d2aa 100644
--- a/airflow-core/src/airflow/models/dagwarning.py
+++ b/airflow-core/src/airflow/models/dagwarning.py
@@ -19,8 +19,9 @@ from __future__ import annotations
from datetime import datetime
from enum import Enum
-from typing import TYPE_CHECKING
+from typing import TYPE_CHECKING, Annotated
+from pydantic import StringConstraints, TypeAdapter
from sqlalchemy import ForeignKeyConstraint, Index, String, Text, delete,
select, true
from sqlalchemy.orm import Mapped, mapped_column, relationship
@@ -34,6 +35,8 @@ from airflow.utils.sqlalchemy import UtcDateTime,
get_dialect_name
if TYPE_CHECKING:
from sqlalchemy.orm import Session
+WARNING_TYPE_MAX_LENGTH = 50
+
class DagWarning(Base):
"""
@@ -45,7 +48,7 @@ class DagWarning(Base):
"""
dag_id: Mapped[str] = mapped_column(StringID(), primary_key=True)
- warning_type: Mapped[str] = mapped_column(String(50), primary_key=True)
+ warning_type: Mapped[str] = mapped_column(String(WARNING_TYPE_MAX_LENGTH),
primary_key=True)
message: Mapped[str] = mapped_column(Text, nullable=False)
timestamp: Mapped[datetime] = mapped_column(UtcDateTime, nullable=False,
default=timezone.utcnow)
@@ -65,7 +68,7 @@ class DagWarning(Base):
def __init__(self, dag_id: str, warning_type: str, message: str, **kwargs):
super().__init__(**kwargs)
self.dag_id = dag_id
- self.warning_type = DagWarningType(warning_type).value # make sure
valid type
+ self.warning_type = get_warning_type_value(warning_type)
self.message = message
def __eq__(self, other) -> bool:
@@ -105,3 +108,25 @@ class DagWarningType(str, Enum):
DUPLICATE_DAG_ID = "duplicate dag id"
NONEXISTENT_POOL = "non-existent pool"
RUNTIME_VARYING_VALUE = "runtime varying value"
+
+
+ImporterWarningType = Annotated[
+ str,
+ StringConstraints(pattern=r"^[a-z][a-z0-9_]*:[a-z0-9_.\-]+$",
max_length=WARNING_TYPE_MAX_LENGTH),
+]
+"""Warning type reported by a Dag importer, prefixed with its namespace (e.g.
``yaml:deprecated_field``)."""
+
+DagWarningTypeValue = DagWarningType | ImporterWarningType
+"""Any valid ``warning_type``: a built-in :class:`DagWarningType` or an
:data:`ImporterWarningType`."""
+
+_warning_type_adapter: TypeAdapter[DagWarningTypeValue] =
TypeAdapter(DagWarningTypeValue)
+
+
+def get_warning_type_value(warning_type: str) -> str:
+ """
+ Return the value to store for ``warning_type``.
+
+ :raises ValueError: if it is neither a :class:`DagWarningType` nor an
:data:`ImporterWarningType`.
+ """
+ validated = _warning_type_adapter.validate_python(warning_type)
+ return validated.value if isinstance(validated, DagWarningType) else
validated
diff --git a/airflow-core/src/airflow/models/errors.py
b/airflow-core/src/airflow/models/errors.py
index 7d967386f79..47fc9811323 100644
--- a/airflow-core/src/airflow/models/errors.py
+++ b/airflow-core/src/airflow/models/errors.py
@@ -18,6 +18,7 @@
from __future__ import annotations
from datetime import datetime
+from pathlib import Path
from sqlalchemy import Integer, String, Text
from sqlalchemy.orm import Mapped, mapped_column
@@ -34,12 +35,17 @@ class ParseImportError(Base):
id: Mapped[int] = mapped_column(Integer, primary_key=True)
timestamp: Mapped[datetime | None] = mapped_column(UtcDateTime,
nullable=True)
filename: Mapped[str | None] = mapped_column(String(1024), nullable=True)
+ source_reference: Mapped[str | None] = mapped_column(String(2000),
nullable=True)
bundle_name: Mapped[str | None] = mapped_column(StringID(), nullable=True)
stacktrace: Mapped[str | None] = mapped_column(Text, nullable=True)
def full_file_path(self) -> str:
"""Return the full file path of the dag."""
- if self.bundle_name is None or self.filename is None:
- raise ValueError("bundle_name and filename must not be None")
+ ref = self.source_reference or self.filename
+ if self.bundle_name is None or ref is None:
+ raise ValueError("bundle_name and (source_reference or filename)
must not be None")
+ ref_path = Path(ref)
+ if ref_path.is_absolute():
+ return ref
bundle = DagBundlesManager().get_bundle(self.bundle_name)
- return "/".join([str(bundle.path), self.filename])
+ return str(bundle.path / ref_path)
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
index c6c5d28aff6..574871897ab 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -2,7 +2,7 @@
import { UseQueryResult } from "@tanstack/react-query";
import { AssetService, AssetStateStoreService, AuthLinksService,
BackfillService, CalendarService, ConfigService, ConnectionService,
DagBundleService, DagParsingService, DagRunService, DagService,
DagSourceService, DagStatsService, DagVersionService, DagWarningService,
DashboardService, DeadlinesService, DependenciesService, EventLogService,
ExperimentalService, ExtraLinksService, GanttService, GridService,
ImportErrorService, JobService, LoginService, MonitorService,
PartitionedDagRunSe [...]
-import { DagRunState, DagSchedulingState, DagWarningType, ReprocessBehavior }
from "../requests/types.gen";
+import { DagRunState, DagSchedulingState, ReprocessBehavior } from
"../requests/types.gen";
export type AssetServiceGetAssetsDefaultResponse = Awaited<ReturnType<typeof
AssetService.getAssets>>;
export type AssetServiceGetAssetsQueryResult<TData =
AssetServiceGetAssetsDefaultResponse, TError = unknown> = UseQueryResult<TData,
TError>;
export const useAssetServiceGetAssetsKey = "AssetServiceGetAssets";
@@ -345,7 +345,7 @@ export const UseDagWarningServiceListDagWarningsKeyFn = ({
dagId, limit, offset,
limit?: number;
offset?: number;
orderBy?: string[];
- warningType?: DagWarningType;
+ warningType?: string;
} = {}, queryKey?: Array<unknown>) => [useDagWarningServiceListDagWarningsKey,
...(queryKey ?? [{ dagId, limit, offset, orderBy, warningType }])];
export type DagServiceGetDagsDefaultResponse = Awaited<ReturnType<typeof
DagService.getDags>>;
export type DagServiceGetDagsQueryResult<TData =
DagServiceGetDagsDefaultResponse, TError = unknown> = UseQueryResult<TData,
TError>;
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
index e1c55f73cf5..01983ce4e99 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
@@ -2,7 +2,7 @@
import { type QueryClient } from "@tanstack/react-query";
import { AssetService, AssetStateStoreService, AuthLinksService,
BackfillService, CalendarService, ConfigService, ConnectionService,
DagBundleService, DagRunService, DagService, DagSourceService, DagStatsService,
DagVersionService, DagWarningService, DashboardService, DeadlinesService,
DependenciesService, EventLogService, ExperimentalService, ExtraLinksService,
GanttService, GridService, ImportErrorService, JobService, LoginService,
MonitorService, PartitionedDagRunService, PluginServic [...]
-import { DagRunState, DagSchedulingState, DagWarningType, ReprocessBehavior }
from "../requests/types.gen";
+import { DagRunState, DagSchedulingState, ReprocessBehavior } from
"../requests/types.gen";
import * as Common from "./common";
/**
* Get Assets
@@ -685,7 +685,7 @@ export const ensureUseDagWarningServiceListDagWarningsData
= (queryClient: Query
limit?: number;
offset?: number;
orderBy?: string[];
- warningType?: DagWarningType;
+ warningType?: string;
} = {}) => queryClient.ensureQueryData({ queryKey:
Common.UseDagWarningServiceListDagWarningsKeyFn({ dagId, limit, offset,
orderBy, warningType }), queryFn: () => DagWarningService.listDagWarnings({
dagId, limit, offset, orderBy, warningType }) });
/**
* Get Dags
@@ -1523,7 +1523,7 @@ export const
ensureUseImportErrorServiceGetImportErrorData = (queryClient: Query
* @param data The data for the request.
* @param data.limit
* @param data.offset
-* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, bundle_name, stacktrace, import_error_id`
+* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, source_reference, bundle_name, stacktrace, import_error_id`
* @param data.filenamePattern Case-insensitive substring match (SQL `ILIKE`).
Slower than `filename_prefix_pattern` on large tables — see "Filtering with
pattern parameters".
* @param data.filenamePrefixPattern Case-sensitive, index-friendly prefix
match. See "Filtering with pattern parameters".
* @param data.filename Exact filename match. Returns only the import error for
this specific file path.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
index 613336b861b..501689dc55d 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
@@ -2,7 +2,7 @@
import { type QueryClient } from "@tanstack/react-query";
import { AssetService, AssetStateStoreService, AuthLinksService,
BackfillService, CalendarService, ConfigService, ConnectionService,
DagBundleService, DagRunService, DagService, DagSourceService, DagStatsService,
DagVersionService, DagWarningService, DashboardService, DeadlinesService,
DependenciesService, EventLogService, ExperimentalService, ExtraLinksService,
GanttService, GridService, ImportErrorService, JobService, LoginService,
MonitorService, PartitionedDagRunService, PluginServic [...]
-import { DagRunState, DagSchedulingState, DagWarningType, ReprocessBehavior }
from "../requests/types.gen";
+import { DagRunState, DagSchedulingState, ReprocessBehavior } from
"../requests/types.gen";
import * as Common from "./common";
/**
* Get Assets
@@ -685,7 +685,7 @@ export const prefetchUseDagWarningServiceListDagWarnings =
(queryClient: QueryCl
limit?: number;
offset?: number;
orderBy?: string[];
- warningType?: DagWarningType;
+ warningType?: string;
} = {}) => queryClient.prefetchQuery({ queryKey:
Common.UseDagWarningServiceListDagWarningsKeyFn({ dagId, limit, offset,
orderBy, warningType }), queryFn: () => DagWarningService.listDagWarnings({
dagId, limit, offset, orderBy, warningType }) });
/**
* Get Dags
@@ -1523,7 +1523,7 @@ export const prefetchUseImportErrorServiceGetImportError
= (queryClient: QueryCl
* @param data The data for the request.
* @param data.limit
* @param data.offset
-* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, bundle_name, stacktrace, import_error_id`
+* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, source_reference, bundle_name, stacktrace, import_error_id`
* @param data.filenamePattern Case-insensitive substring match (SQL `ILIKE`).
Slower than `filename_prefix_pattern` on large tables — see "Filtering with
pattern parameters".
* @param data.filenamePrefixPattern Case-sensitive, index-friendly prefix
match. See "Filtering with pattern parameters".
* @param data.filename Exact filename match. Returns only the import error for
this specific file path.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
index aaa024a638c..4951fc96970 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -2,7 +2,7 @@
import { UseMutationOptions, UseQueryOptions, useMutation, useQuery } from
"@tanstack/react-query";
import { AssetService, AssetStateStoreService, AuthLinksService,
BackfillService, CalendarService, ConfigService, ConnectionService,
DagBundleService, DagParsingService, DagRunService, DagService,
DagSourceService, DagStatsService, DagVersionService, DagWarningService,
DashboardService, DeadlinesService, DependenciesService, EventLogService,
ExperimentalService, ExtraLinksService, GanttService, GridService,
ImportErrorService, JobService, LoginService, MonitorService,
PartitionedDagRunSe [...]
-import { AssetStateStoreBody, BackfillPostBody, BulkBody_BulkDAGBody_,
BulkBody_BulkDAGRunBody_, BulkBody_BulkTaskInstanceBody_,
BulkBody_ConnectionBody_, BulkBody_PoolBody_, BulkBody_VariableBody_,
BulkDAGRunClearBody, ClearPartitionsBody, ClearTaskInstancesBody,
ConnectionBody, ConnectionTestRequestBody, CreateAssetEventsBody, DAGPatchBody,
DAGRunClearBody, DAGRunPatchBody, DAGRunsBatchBody, DagRunState,
DagSchedulingState, DagWarningType, GenerateTokenBody, MaterializeAssetBody,
Patch [...]
+import { AssetStateStoreBody, BackfillPostBody, BulkBody_BulkDAGBody_,
BulkBody_BulkDAGRunBody_, BulkBody_BulkTaskInstanceBody_,
BulkBody_ConnectionBody_, BulkBody_PoolBody_, BulkBody_VariableBody_,
BulkDAGRunClearBody, ClearPartitionsBody, ClearTaskInstancesBody,
ConnectionBody, ConnectionTestRequestBody, CreateAssetEventsBody, DAGPatchBody,
DAGRunClearBody, DAGRunPatchBody, DAGRunsBatchBody, DagRunState,
DagSchedulingState, GenerateTokenBody, MaterializeAssetBody,
PatchTaskInstanceBody [...]
import * as Common from "./common";
/**
* Get Assets
@@ -685,7 +685,7 @@ export const useDagWarningServiceListDagWarnings = <TData =
Common.DagWarningSer
limit?: number;
offset?: number;
orderBy?: string[];
- warningType?: DagWarningType;
+ warningType?: string;
} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseDagWarningServiceListDagWarningsKeyFn({ dagId, limit, offset,
orderBy, warningType }, queryKey), queryFn: () =>
DagWarningService.listDagWarnings({ dagId, limit, offset, orderBy, warningType
}) as TData, ...options });
/**
* Get Dags
@@ -1523,7 +1523,7 @@ export const useImportErrorServiceGetImportError = <TData
= Common.ImportErrorSe
* @param data The data for the request.
* @param data.limit
* @param data.offset
-* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, bundle_name, stacktrace, import_error_id`
+* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, source_reference, bundle_name, stacktrace, import_error_id`
* @param data.filenamePattern Case-insensitive substring match (SQL `ILIKE`).
Slower than `filename_prefix_pattern` on large tables — see "Filtering with
pattern parameters".
* @param data.filenamePrefixPattern Case-sensitive, index-friendly prefix
match. See "Filtering with pattern parameters".
* @param data.filename Exact filename match. Returns only the import error for
this specific file path.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
index 9add7eba93e..6c91b91e143 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
@@ -2,7 +2,7 @@
import { UseQueryOptions, useSuspenseQuery } from "@tanstack/react-query";
import { AssetService, AssetStateStoreService, AuthLinksService,
BackfillService, CalendarService, ConfigService, ConnectionService,
DagBundleService, DagRunService, DagService, DagSourceService, DagStatsService,
DagVersionService, DagWarningService, DashboardService, DeadlinesService,
DependenciesService, EventLogService, ExperimentalService, ExtraLinksService,
GanttService, GridService, ImportErrorService, JobService, LoginService,
MonitorService, PartitionedDagRunService, PluginServic [...]
-import { DagRunState, DagSchedulingState, DagWarningType, ReprocessBehavior }
from "../requests/types.gen";
+import { DagRunState, DagSchedulingState, ReprocessBehavior } from
"../requests/types.gen";
import * as Common from "./common";
/**
* Get Assets
@@ -685,7 +685,7 @@ export const useDagWarningServiceListDagWarningsSuspense =
<TData = Common.DagWa
limit?: number;
offset?: number;
orderBy?: string[];
- warningType?: DagWarningType;
+ warningType?: string;
} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseDagWarningServiceListDagWarningsKeyFn({ dagId, limit, offset,
orderBy, warningType }, queryKey), queryFn: () =>
DagWarningService.listDagWarnings({ dagId, limit, offset, orderBy, warningType
}) as TData, ...options });
/**
* Get Dags
@@ -1523,7 +1523,7 @@ export const useImportErrorServiceGetImportErrorSuspense
= <TData = Common.Impor
* @param data The data for the request.
* @param data.limit
* @param data.offset
-* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, bundle_name, stacktrace, import_error_id`
+* @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, source_reference, bundle_name, stacktrace, import_error_id`
* @param data.filenamePattern Case-insensitive substring match (SQL `ILIKE`).
Slower than `filename_prefix_pattern` on large tables — see "Filtering with
pattern parameters".
* @param data.filenamePrefixPattern Case-sensitive, index-friendly prefix
match. See "Filtering with pattern parameters".
* @param data.filename Exact filename match. Returns only the import error for
this specific file path.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 11bd522c626..b028696f3a8 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -4513,6 +4513,17 @@ export const $DAGSourceResponse = {
dag_display_name: {
type: 'string',
title: 'Dag Display Name'
+ },
+ language: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Language'
}
},
type: 'object',
@@ -4588,7 +4599,17 @@ export const $DAGWarningResponse = {
title: 'Dag Id'
},
warning_type: {
- '$ref': '#/components/schemas/DagWarningType'
+ anyOf: [
+ {
+ '$ref': '#/components/schemas/DagWarningType'
+ },
+ {
+ type: 'string',
+ maxLength: 50,
+ pattern: '^[a-z][a-z0-9_]*:[a-z0-9_.\\-]+$'
+ }
+ ],
+ title: 'Warning Type'
},
message: {
type: 'string',
@@ -6130,6 +6151,17 @@ export const $ImportErrorResponse = {
type: 'string',
title: 'Filename'
},
+ source_reference: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Source Reference'
+ },
bundle_name: {
anyOf: [
{
@@ -6153,7 +6185,7 @@ export const $ImportErrorResponse = {
}
},
type: 'object',
- required: ['import_error_id', 'timestamp', 'filename', 'bundle_name',
'stack_trace', 'file_token'],
+ required: ['import_error_id', 'timestamp', 'filename', 'source_reference',
'bundle_name', 'stack_trace', 'file_token'],
title: 'ImportErrorResponse',
description: 'Import Error Response.'
} as const;
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
index f796d225c6e..d0a1a936ed9 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts
@@ -3588,7 +3588,7 @@ export class ImportErrorService {
* @param data The data for the request.
* @param data.limit
* @param data.offset
- * @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, bundle_name, stacktrace, import_error_id`
+ * @param data.orderBy Attributes to order by, multi criteria sort is
supported. Prefix with `-` for descending order. Supported attributes: `id,
timestamp, filename, source_reference, bundle_name, stacktrace, import_error_id`
* @param data.filenamePattern Case-insensitive substring match (SQL
`ILIKE`). Slower than `filename_prefix_pattern` on large tables — see
"Filtering with pattern parameters".
* @param data.filenamePrefixPattern Case-sensitive, index-friendly prefix
match. See "Filtering with pattern parameters".
* @param data.filename Exact filename match. Returns only the import
error for this specific file path.
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index 8f70583bf01..ceb90161cee 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -1163,6 +1163,7 @@ export type DAGSourceResponse = {
dag_id: string;
version_number: number | null;
dag_display_name: string;
+ language?: string | null;
};
/**
@@ -1194,7 +1195,7 @@ export type DAGWarningCollectionResponse = {
*/
export type DAGWarningResponse = {
dag_id: string;
- warning_type: DagWarningType;
+ warning_type: DagWarningType | string;
message: string;
timestamp: string;
dag_display_name: string;
@@ -1662,6 +1663,7 @@ export type ImportErrorResponse = {
import_error_id: number;
timestamp: string;
filename: string;
+ source_reference: string | null;
bundle_name: string | null;
stack_trace: string;
/**
@@ -3718,7 +3720,7 @@ export type ListDagWarningsData = {
* Attributes to order by, multi criteria sort is supported. Prefix with
`-` for descending order. Supported attributes: `dag_id, warning_type, message,
timestamp`
*/
orderBy?: Array<(string)>;
- warningType?: DagWarningType | null;
+ warningType?: DagWarningType | string | null;
};
export type ListDagWarningsResponse = DAGWarningCollectionResponse;
@@ -4520,7 +4522,7 @@ export type GetImportErrorsData = {
limit?: number;
offset?: number;
/**
- * Attributes to order by, multi criteria sort is supported. Prefix with
`-` for descending order. Supported attributes: `id, timestamp, filename,
bundle_name, stacktrace, import_error_id`
+ * Attributes to order by, multi criteria sort is supported. Prefix with
`-` for descending order. Supported attributes: `id, timestamp, filename,
source_reference, bundle_name, stacktrace, import_error_id`
*/
orderBy?: Array<(string)>;
};
diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/Code/Code.tsx
b/airflow-core/src/airflow/ui/src/pages/Dag/Code/Code.tsx
index ca274462779..b8874d6b262 100644
--- a/airflow-core/src/airflow/ui/src/pages/Dag/Code/Code.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Dag/Code/Code.tsx
@@ -161,6 +161,8 @@ export const Code = () => {
? translate("code.noCode")
: (compareCode?.content ?? "");
+ const language: string = code?.language ?? "python";
+
const codeStatus = (
<>
<ErrorAlert
@@ -183,7 +185,11 @@ export const Code = () => {
<FileLocation fileloc={dag.fileloc}
relativeFileloc={dag.relative_fileloc} />
)}
<Box flex={1} minH={0}>
- <CodeDiffViewer modifiedCode={displayedCode}
originalCode={displayedCompareCode} />
+ <CodeDiffViewer
+ language={language}
+ modifiedCode={displayedCode}
+ originalCode={displayedCompareCode}
+ />
</Box>
</Box>
) : (
@@ -206,7 +212,7 @@ export const Code = () => {
<Box flex={1} minH={0}>
<Editor
beforeMount={beforeMount}
- language="python"
+ language={language}
options={editorOptions}
theme={theme}
value={displayedCode}
diff --git a/airflow-core/src/airflow/utils/db.py
b/airflow-core/src/airflow/utils/db.py
index 950af812a7f..9639594f721 100644
--- a/airflow-core/src/airflow/utils/db.py
+++ b/airflow-core/src/airflow/utils/db.py
@@ -117,7 +117,7 @@ _REVISION_HEADS_MAP: dict[str, str] = {
"3.1.8": "509b94a1042d",
"3.2.0": "1d6611b6ab7c",
"3.3.0": "d2f4e1b3c5a7",
- "3.4.0": "a61f0c9d2b47",
+ "3.4.0": "e5a91c7f42b3",
}
# Prefix used to identify tables holding data moved during migration.
diff --git a/airflow-core/tests/unit/api/common/test_delete_dag.py
b/airflow-core/tests/unit/api/common/test_delete_dag.py
index eea327bdfcd..14f96e3b6fe 100644
--- a/airflow-core/tests/unit/api/common/test_delete_dag.py
+++ b/airflow-core/tests/unit/api/common/test_delete_dag.py
@@ -24,8 +24,11 @@ from sqlalchemy import func, select
from airflow.api.common.delete_dag import delete_dag
from airflow.models import DagModel
+from airflow.models.errors import ParseImportError
from airflow.providers.standard.operators.empty import EmptyOperator
+from tests_common.test_utils.db import clear_db_import_errors
+
if TYPE_CHECKING:
from airflow.serialization.definitions.dag import SerializedDAG
@@ -36,6 +39,13 @@ pytestmark = [pytest.mark.db_test,
pytest.mark.need_serialized_dag]
DAG_ID = "dag_to_delete"
[email protected]
+def clean_import_errors():
+ clear_db_import_errors()
+ yield
+ clear_db_import_errors()
+
+
def test_delete_dag_does_not_read_back_deleted_row_keys(dag_maker:
DagMaker[SerializedDAG], session):
"""
delete_dag must not ask the database for the keys of the rows it deletes.
@@ -76,3 +86,42 @@ def
test_delete_dag_does_not_read_back_deleted_row_keys(dag_maker: DagMaker[Seri
)
assert
session.scalar(select(func.count()).select_from(DagModel).where(DagModel.dag_id
== DAG_ID)) == 0
+
+
[email protected](
+ "matching_column",
+ [
+ pytest.param("source_reference", id="by-source-reference"),
+ pytest.param("filename", id="by-relative-filename"),
+ ],
+)
[email protected]("clean_import_errors")
+def test_delete_dag_removes_import_errors_for_its_file(
+ dag_maker: DagMaker[SerializedDAG], session, matching_column
+):
+ with dag_maker(DAG_ID, session=session):
+ EmptyOperator(task_id="task")
+ session.commit()
+ dag_model = session.scalar(select(DagModel).where(DagModel.dag_id ==
DAG_ID))
+
+ session.add_all(
+ [
+ ParseImportError(
+ bundle_name=dag_model.bundle_name,
+ stacktrace="matching",
+ **{matching_column: dag_model.relative_fileloc},
+ ),
+ ParseImportError(
+ bundle_name=dag_model.bundle_name,
+ filename="unrelated.py",
+ source_reference="unrelated.py",
+ stacktrace="unrelated",
+ ),
+ ]
+ )
+ session.commit()
+
+ delete_dag(DAG_ID, session=session)
+ session.commit()
+
+ assert session.scalars(select(ParseImportError.stacktrace)).all() ==
["unrelated"]
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_sources.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_sources.py
index 93083102fed..6c07e4d3892 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_sources.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_sources.py
@@ -122,6 +122,7 @@ class TestGetDAGSource:
"dag_id": TEST_DAG_ID,
"version_number": 1,
"dag_display_name": TEST_DAG_DISPLAY_NAME,
+ "language": "python",
}
assert response.headers["Content-Type"].startswith("application/json")
@@ -167,6 +168,7 @@ class TestGetDAGSource:
"dag_id": TEST_DAG_ID,
"version_number": 2,
"dag_display_name": TEST_DAG_DISPLAY_NAME,
+ "language": "python",
}
def test_should_respond_406_unsupport_mime_type(self, test_client,
test_dag):
@@ -225,6 +227,7 @@ class TestGetDAGSource:
"dag_id": TEST_DAG_ID,
"version_number": 1,
"dag_display_name": TEST_DAG_DISPLAY_NAME,
+ "language": "python",
}
mock_get_auth_manager.return_value.get_authorized_dag_ids.assert_called_once_with(user=mock.ANY)
@@ -248,4 +251,5 @@ class TestGetDAGSource:
"dag_id": TEST_DAG_ID,
"version_number": 1,
"dag_display_name": TEST_DAG_DISPLAY_NAME,
+ "language": "python",
}
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_warning.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_warning.py
index ade87cccf60..7a15de834c9 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_warning.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_warning.py
@@ -119,3 +119,22 @@ class TestGetDagWarnings:
response_json["detail"][0]["msg"]
== "Input should be 'asset conflict', 'duplicate dag id',
'non-existent pool' or 'runtime varying value'"
)
+
+ def test_get_dag_warnings_non_matching_type(self, test_client):
+ response = test_client.get("/dagWarnings", params={"warning_type":
"yaml:unused_type"})
+ response_json = response.json()
+ assert response.status_code == 200
+ assert response_json["total_entries"] == 0
+ assert response_json["dag_warnings"] == []
+
+ def test_get_dag_warnings_importer_defined_type(self, test_client,
session):
+ session.add(DagWarning(DAG1_ID, "yaml:schema_violation", "importer
message"))
+ session.commit()
+
+ response = test_client.get("/dagWarnings", params={"warning_type":
"yaml:schema_violation"})
+
+ assert response.status_code == 200
+ response_json = response.json()
+ assert response_json["total_entries"] == 1
+ assert response_json["dag_warnings"][0]["warning_type"] ==
"yaml:schema_violation"
+ assert response_json["dag_warnings"][0]["message"] == "importer
message"
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
index 1bde7d2e1a3..02f5878db2a 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_import_error.py
@@ -21,7 +21,7 @@ from typing import TYPE_CHECKING
from unittest import mock
import pytest
-from sqlalchemy import delete
+from sqlalchemy import delete, update
from airflow.api_fastapi.auth.managers.models.resource_details import
AccessView, DagDetails
from airflow.api_fastapi.core_api.routes.public.import_error import
REDACTED_STACKTRACE
@@ -171,6 +171,7 @@ class TestGetImportError:
{
"timestamp": from_datetime_to_zulu_without_ms(TIMESTAMP1),
"filename": FILENAME1,
+ "source_reference": None,
"stack_trace": STACKTRACE1,
"bundle_name": BUNDLE_NAME,
},
@@ -181,6 +182,7 @@ class TestGetImportError:
{
"timestamp": from_datetime_to_zulu_without_ms(TIMESTAMP2),
"filename": FILENAME2,
+ "source_reference": None,
"stack_trace": STACKTRACE2,
"bundle_name": BUNDLE_NAME,
},
@@ -191,6 +193,7 @@ class TestGetImportError:
{
"timestamp": from_datetime_to_zulu_without_ms(TIMESTAMP3),
"filename": FILENAME3,
+ "source_reference": None,
"stack_trace": STACKTRACE3,
"bundle_name": BUNDLE_NAME,
},
@@ -230,6 +233,35 @@ class TestGetImportError:
)
assert response.json() == expected_body
+
@mock.patch("airflow.api_fastapi.core_api.routes.public.import_error.get_auth_manager")
+ def test_get_import_error_with_source_reference(
+ self, mock_get_auth_manager, test_client, session,
permitted_dag_model_all, url_safe_serializer
+ ):
+ error = ParseImportError(
+ bundle_name=BUNDLE_NAME,
+ filename=FILENAME1,
+ source_reference="archive.zip/dags/my_dag.py",
+ stacktrace=STACKTRACE1,
+ timestamp=TIMESTAMP1,
+ )
+ session.add(error)
+ session.commit()
+
+ set_mock_auth_manager__get_authorized_dag_ids(mock_get_auth_manager,
permitted_dag_model_all)
+ response = test_client.get(f"/importErrors/{error.id}")
+ assert response.status_code == 200
+ assert response.json() == {
+ "import_error_id": error.id,
+ "timestamp": from_datetime_to_zulu_without_ms(TIMESTAMP1),
+ "filename": FILENAME1,
+ "source_reference": "archive.zip/dags/my_dag.py",
+ "stack_trace": STACKTRACE1,
+ "bundle_name": BUNDLE_NAME,
+ "file_token": url_safe_serializer.dumps(
+ {"bundle_name": BUNDLE_NAME, "relative_fileloc": FILENAME1}
+ ),
+ }
+
def test_should_raises_401_unauthenticated(self,
unauthenticated_test_client, import_errors):
import_error_id = import_errors[0].id
response =
unauthenticated_test_client.get(f"/importErrors/{import_error_id}")
@@ -275,6 +307,7 @@ class TestGetImportError:
"import_error_id": import_error_id,
"timestamp": from_datetime_to_zulu_without_ms(TIMESTAMP1),
"filename": FILENAME1,
+ "source_reference": None,
"stack_trace": "REDACTED - you do not have read permission on all
Dags in the file",
"bundle_name": BUNDLE_NAME,
"file_token": url_safe_serializer.dumps(
@@ -316,6 +349,7 @@ class TestGetImportError:
"import_error_id": import_error_id,
"timestamp": from_datetime_to_zulu_without_ms(TIMESTAMP1),
"filename": FILENAME1,
+ "source_reference": None,
"stack_trace": STACKTRACE1,
"bundle_name": BUNDLE_NAME,
"file_token": url_safe_serializer.dumps(
@@ -643,6 +677,45 @@ class TestGetImportErrors:
access_view=AccessView.IMPORT_ERRORS_ALL, user=mock.ANY,
team_name=None
)
+ @pytest.mark.parametrize(
+ ("order_by", "expected_filenames"),
+ [
+ ("source_reference", [FILENAME2, FILENAME3, FILENAME1]),
+ ("-source_reference", [FILENAME1, FILENAME3, FILENAME2]),
+ ],
+ )
+
@mock.patch("airflow.api_fastapi.core_api.routes.public.import_error.get_auth_manager")
+ def test_get_import_errors_source_reference(
+ self,
+ mock_get_auth_manager,
+ test_client,
+ session,
+ order_by,
+ expected_filenames,
+ permitted_dag_model_all,
+ ):
+ source_references = {
+ FILENAME1: f"c.zip/{FILENAME1}",
+ FILENAME2: f"a.zip/{FILENAME2}",
+ FILENAME3: f"b.zip/{FILENAME3}",
+ }
+ for filename, source_reference in source_references.items():
+ session.execute(
+ update(ParseImportError)
+ .where(ParseImportError.filename == filename)
+ .values(source_reference=source_reference)
+ )
+ session.commit()
+ set_mock_auth_manager__get_authorized_dag_ids(mock_get_auth_manager,
permitted_dag_model_all)
+ set_mock_auth_manager__batch_is_authorized_dag(mock_get_auth_manager,
True)
+
+ response = test_client.get("/importErrors", params={"order_by":
order_by})
+
+ assert response.status_code == 200
+ assert [(e["filename"], e["source_reference"]) for e in
response.json()["import_errors"]] == [
+ (filename, source_references[filename]) for filename in
expected_filenames
+ ]
+
def test_should_raises_401_unauthenticated(self,
unauthenticated_test_client):
response = unauthenticated_test_client.get("/importErrors")
assert response.status_code == 401
@@ -785,6 +858,7 @@ class TestGetImportErrors:
"import_error_id": import_errors[0].id,
"timestamp": from_datetime_to_zulu_without_ms(TIMESTAMP1),
"filename": FILENAME1,
+ "source_reference": None,
"stack_trace": expected_stack_trace,
"bundle_name": BUNDLE_NAME,
"file_token": url_safe_serializer.dumps(
diff --git a/airflow-core/tests/unit/models/test_dagcode.py
b/airflow-core/tests/unit/models/test_dagcode.py
index 1a0e16d95e5..419a8f0cbd5 100644
--- a/airflow-core/tests/unit/models/test_dagcode.py
+++ b/airflow-core/tests/unit/models/test_dagcode.py
@@ -262,3 +262,11 @@ class TestDagCode:
refreshed = DagCode.get_latest_dagcode(dag.dag_id)
assert refreshed.fileloc == dag.fileloc
assert refreshed.source_code_hash == original_hash
+
+ def test_language_defaults_to_python(self, dag_maker):
+ """Code written without an explicit language is recorded as Python."""
+ with dag_maker("dag_language_default") as dag:
+ pass
+ sync_dag_to_db(dag)
+
+ assert DagCode.get_latest_dagcode(dag.dag_id).language == "python"
diff --git a/airflow-core/tests/unit/models/test_dagwarning.py
b/airflow-core/tests/unit/models/test_dagwarning.py
index 0c5a3d3f969..0f07d7a9192 100644
--- a/airflow-core/tests/unit/models/test_dagwarning.py
+++ b/airflow-core/tests/unit/models/test_dagwarning.py
@@ -25,13 +25,45 @@ from sqlalchemy import select
from sqlalchemy.exc import OperationalError
from airflow.models import DagModel
-from airflow.models.dagwarning import DagWarning
+from airflow.models.dagwarning import DagWarning, DagWarningType,
get_warning_type_value
from tests_common.test_utils.db import clear_db_dags
-pytestmark = pytest.mark.db_test
-
+class TestGetWarningTypeValue:
+ @pytest.mark.parametrize(
+ ("warning_type", "expected"),
+ [
+ pytest.param(DagWarningType.NONEXISTENT_POOL, "non-existent pool",
id="enum-member"),
+ pytest.param("non-existent pool", "non-existent pool",
id="built-in-string"),
+ pytest.param("yaml:schema_violation", "yaml:schema_violation",
id="importer-type"),
+ pytest.param(f"yaml:{'x' * 45}", f"yaml:{'x' * 45}",
id="importer-type-at-column-length"),
+ ],
+ )
+ def test_valid_type_returns_stored_value(self, warning_type, expected):
+ assert get_warning_type_value(warning_type) == expected
+
+ @pytest.mark.parametrize(
+ "warning_type",
+ [
+ pytest.param(1, id="non-string"),
+ pytest.param("yaml schema violation", id="no-namespace"),
+ pytest.param("non existent pool", id="mistyped-built-in-type"),
+ pytest.param("YAML:Schema_Violation", id="uppercase"),
+ pytest.param(":schema_violation", id="empty-namespace"),
+ pytest.param("yaml:", id="empty-type"),
+ pytest.param(f"yaml:{'x' * 46}", id="longer-than-column"),
+ ],
+ )
+ def test_invalid_type_is_rejected(self, warning_type):
+ with pytest.raises(ValueError, match="validation error"):
+ get_warning_type_value(warning_type)
+
+ def test_dag_warning_stores_importer_type(self):
+ assert DagWarning("dag_1", "yaml:schema_violation",
"message").warning_type == "yaml:schema_violation"
+
+
[email protected]_test
class TestDagWarning:
def setup_method(self):
clear_db_dags()
diff --git a/airflow-core/tests/unit/models/test_errors.py
b/airflow-core/tests/unit/models/test_errors.py
new file mode 100644
index 00000000000..da5d1f36c27
--- /dev/null
+++ b/airflow-core/tests/unit/models/test_errors.py
@@ -0,0 +1,86 @@
+# 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.
+from __future__ import annotations
+
+from pathlib import Path
+from unittest import mock
+
+import pytest
+
+from airflow.models.errors import ParseImportError
+
+BUNDLE_PATH = Path("/bundle")
+
+
[email protected]
+def bundle_path():
+ with mock.patch("airflow.models.errors.DagBundlesManager", autospec=True)
as manager:
+ manager.return_value.get_bundle.return_value.path = BUNDLE_PATH
+ yield
+
+
+class TestFullFilePath:
+ @pytest.mark.usefixtures("bundle_path")
+ @pytest.mark.parametrize(
+ ("source_reference", "expected"),
+ [
+ pytest.param("dags/my_dag.py", "/bundle/dags/my_dag.py",
id="relative-file"),
+ pytest.param(
+ "archive.zip/dags/my_dag.py",
+ "/bundle/archive.zip/dags/my_dag.py",
+ id="relative-archive-member",
+ ),
+ pytest.param(
+ "/elsewhere/archive.zip/dags/my_dag.py",
+ "/elsewhere/archive.zip/dags/my_dag.py",
+ id="absolute-archive-member",
+ ),
+ pytest.param("/bundle/dags/my_dag.py", "/bundle/dags/my_dag.py",
id="already-under-bundle"),
+ pytest.param(
+ "/bundle-sibling/dags/my_dag.py",
+ "/bundle-sibling/dags/my_dag.py",
+ id="sibling-bundle-prefix",
+ ),
+ ],
+ )
+ def test_resolves_reference_against_bundle(self, source_reference,
expected):
+ error = ParseImportError(bundle_name="my-bundle",
source_reference=source_reference)
+
+ assert error.full_file_path() == expected
+
+ @pytest.mark.usefixtures("bundle_path")
+ def test_resolves_filename_when_source_reference_is_none(self):
+ error = ParseImportError(bundle_name="my-bundle",
filename="dags/my_dag.py")
+
+ assert error.full_file_path() == "/bundle/dags/my_dag.py"
+
+ @pytest.mark.parametrize(
+ ("bundle_name", "source_reference", "filename"),
+ [
+ pytest.param(None, "dags/my_dag.py", None, id="missing-bundle"),
+ pytest.param("my-bundle", None, None, id="missing-reference"),
+ ],
+ )
+ def test_raises_when_reference_is_incomplete(self, bundle_name,
source_reference, filename):
+ error = ParseImportError(
+ bundle_name=bundle_name, source_reference=source_reference,
filename=filename
+ )
+
+ with pytest.raises(
+ ValueError, match=r"bundle_name and \(source_reference or
filename\) must not be None"
+ ):
+ error.full_file_path()
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index b59a1625dff..924d3420928 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -489,6 +489,7 @@ class DAGSourceResponse(BaseModel):
dag_id: Annotated[str, Field(title="Dag Id")]
version_number: Annotated[int | None, Field(title="Version Number")]
dag_display_name: Annotated[str, Field(title="Dag Display Name")]
+ language: Annotated[str | None, Field(title="Language")] = None
class DAGTagCollectionResponse(BaseModel):
@@ -500,6 +501,12 @@ class DAGTagCollectionResponse(BaseModel):
total_entries: Annotated[int, Field(title="Total Entries")]
+class WarningType(RootModel[str]):
+ root: Annotated[
+ str, Field(max_length=50, pattern="^[a-z][a-z0-9_]*:[a-z0-9_.\\-]+$",
title="Warning Type")
+ ]
+
+
class DagBundleDetailResponse(BaseModel):
"""
Dag bundle serializer for the single-bundle response.
@@ -925,6 +932,7 @@ class ImportErrorResponse(BaseModel):
import_error_id: Annotated[int, Field(title="Import Error Id")]
timestamp: Annotated[datetime, Field(title="Timestamp")]
filename: Annotated[str, Field(title="Filename")]
+ source_reference: Annotated[str | None, Field(title="Source Reference")]
bundle_name: Annotated[str | None, Field(title="Bundle Name")]
stack_trace: Annotated[str, Field(title="Stack Trace")]
file_token: Annotated[
@@ -2157,7 +2165,7 @@ class DAGWarningResponse(BaseModel):
"""
dag_id: Annotated[str, Field(title="Dag Id")]
- warning_type: DagWarningType
+ warning_type: Annotated[DagWarningType | WarningType, Field(title="Warning
Type")]
message: Annotated[str, Field(title="Message")]
timestamp: Annotated[datetime, Field(title="Timestamp")]
dag_display_name: Annotated[str, Field(title="Dag Display Name")]
diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
index f1eaf2e89ef..24ed44e2620 100644
--- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
+++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
@@ -1151,6 +1151,7 @@ class TestDagOperations:
import_error_id=0,
timestamp=datetime.datetime(2025, 1, 1, 0, 0, 0),
filename="filename",
+ source_reference=None,
bundle_name="bundle_name",
stack_trace="stack_trace",
file_token="file_token",