This is an automated email from the ASF dual-hosted git repository.
kaxil 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 3a944127dd8 Add a Dag bundle detail page listing the files in the
bundle (#73009)
3a944127dd8 is described below
commit 3a944127dd8088415e32ae8009788afba7256708
Author: Kaxil Naik <[email protected]>
AuthorDate: Tue Sep 22 00:34:54 2026 +0100
Add a Dag bundle detail page listing the files in the bundle (#73009)
Add GET /api/v2/dagBundles/{bundle_name} and
GET /api/v2/dagBundles/{bundle_name}/files, and a page behind each row of
the Dag
Bundles list showing the bundle's files with the Dags each one defines,
when it
was last parsed, how long that took, and its import-error count.
The list view answers "has my commit landed" but not "what in it is
broken". The
gap it cannot close is a file that fails before defining a Dag: there is no
DagModel row for it, so no per-Dag view in the UI mentions it, and the
bundle's
import-error count says only that something is wrong. That file is exactly
what
someone opens this page to find.
Files therefore come from two sources. Files holding a Dag the caller may
read
are aggregated from dag, which is also the only record of parse time and
duration Airflow keeps -- a file whose Dags were all removed keeps none of
its
own. Files that recorded an import error without registering a Dag are
added on
the same admin-by-default terms as GET /importErrors, since they have no
Dag to
authorize on. The two halves are merged in the api-server rather than
unioned in
the database: they differ in shape and in authorization, and the set being
merged
is one bundle's files.
Dag rows are never deleted, which decides what the counts can mean. A file
that
fails to import has every Dag in it marked stale by
update_dag_parsing_results_in_db, and a file dropped from the bundle has
the same
done by deactivate_deleted_dags. So dag_count counts live Dags only,
agreeing
with GET /dags, which excludes stale Dags by default -- counting them would
over-report in the one case the page exists to diagnose. A file then stays
listed
on the strength of a live Dag or an import error, rather than of a Dag row:
a
file that has just broken keeps its row instead of vanishing, and a file
deleted
long ago drops out instead of lingering with a count of Dags that no longer
exist.
A file's import-error count opens the stack trace behind it, read through
GET /importErrors rather than returned with the file list: a stack trace is
unbounded and the file list is polled, and that endpoint already decides
who may
read an import error. It renders in DagImportErrorModal, which moves from
pages/Dag to components now that a second page uses it -- components already
reached across into pages for it, which was backwards.
Staleness decides the parse columns too. last_parsed_time is the file's last
known parse, so a broken file still says when it last worked, while
last_parse_duration is taken from live rows only: a Dag removed from a file
keeps
its row frozen at an older, slower parse, and an unrestricted MAX would
pair that
duration with the newer timestamp and report a parse that never happened.
Both routes resolve the bundle through the same filter the collection route
uses,
so a bundle holding no registered Dag is reachable on exactly the terms the
collection grants it, and a bundle the caller may not see is a 404 rather
than a
403 -- a response cannot be used to confirm which bundle names a deployment
has.
Neither the bundle classpath nor its configured kwargs are exposed. The
classpath
would mean reading dag_bundle_config_list in the api-server, which is the
reason
refresh_interval stayed out, and bundle kwargs are free-form and a
plausible home
for credentials.
---
airflow-core/docs/security/api_permissions_ref.rst | 8 +
.../api_fastapi/core_api/datamodels/dag_bundles.py | 38 +++
.../core_api/openapi/v2-rest-api-generated.yaml | 268 +++++++++++++++++++++
.../core_api/routes/public/dag_bundles.py | 254 +++++++++++++++++--
.../src/airflow/ui/openapi-gen/queries/common.ts | 14 ++
.../ui/openapi-gen/queries/ensureQueryData.ts | 40 +++
.../src/airflow/ui/openapi-gen/queries/prefetch.ts | 40 +++
.../src/airflow/ui/openapi-gen/queries/queries.ts | 40 +++
.../src/airflow/ui/openapi-gen/queries/suspense.ts | 40 +++
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 166 +++++++++++++
.../ui/openapi-gen/requests/services.gen.ts | 70 +++++-
.../airflow/ui/openapi-gen/requests/types.gen.ts | 134 +++++++++++
.../airflow/ui/public/i18n/locales/en/browse.json | 18 +-
.../airflow/ui/public/i18n/locales/en/common.json | 2 +
.../airflow/ui/src/components/DagBundleVersion.tsx | 63 +++++
.../ui/src/components/DagDeactivatedBanner.tsx | 2 +-
.../Dag => components}/DagImportErrorModal.tsx | 0
.../src/airflow/ui/src/components/HeaderCard.tsx | 2 +-
.../airflow/ui/src/components/ImportErrorCount.tsx | 43 ++++
.../airflow/ui/src/pages/DagBundle/BundleFiles.tsx | 144 +++++++++++
.../airflow/ui/src/pages/DagBundle/DagBundle.tsx | 62 +++++
.../ui/src/pages/DagBundle/FileImportError.tsx | 78 ++++++
.../src/airflow/ui/src/pages/DagBundle/Header.tsx | 84 +++++++
.../src/airflow/ui/src/pages/DagBundle/index.tsx | 19 ++
.../airflow/ui/src/pages/DagBundles/DagBundles.tsx | 87 ++-----
.../ui/src/queries/useDagBundleRefetchInterval.ts | 38 +++
airflow-core/src/airflow/ui/src/router.tsx | 5 +
.../core_api/routes/public/test_dag_bundles.py | 266 +++++++++++++++++++-
.../src/airflowctl/api/datamodels/generated.py | 98 ++++++++
providers/fab/docs/auth-manager/access-control.rst | 8 +
30 files changed, 2041 insertions(+), 90 deletions(-)
diff --git a/airflow-core/docs/security/api_permissions_ref.rst
b/airflow-core/docs/security/api_permissions_ref.rst
index 6b92ab32240..c930fdb9f07 100644
--- a/airflow-core/docs/security/api_permissions_ref.rst
+++ b/airflow-core/docs/security/api_permissions_ref.rst
@@ -186,6 +186,14 @@ source code so it stays up to date as endpoints are added
or changed.
- ``/api/v2/dagBundles``
- ``DAG``
- ``GET``
+ * - ``GET``
+ - ``/api/v2/dagBundles/{bundle_name}``
+ - ``DAG``
+ - ``GET``
+ * - ``GET``
+ - ``/api/v2/dagBundles/{bundle_name}/files``
+ - ``DAG``
+ - ``GET``
* - ``GET``
- ``/api/v2/dagSources/{dag_id}``
- ``DAG.CODE``
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_bundles.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_bundles.py
index 198dd2898c9..7429f8af1d9 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_bundles.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_bundles.py
@@ -64,3 +64,41 @@ class DagBundleCollectionResponse(BaseModel):
dag_bundles: Iterable[DagBundleResponse]
total_entries: int
+
+
+class DagBundleDetailResponse(DagBundleResponse):
+ """Dag bundle serializer for the single-bundle response."""
+
+ dag_count: int = Field(
+ description=(
+ "Number of live Dags recorded against the bundle that the caller
is permitted to see, "
+ "counted on the same terms as ``GET /dags``."
+ )
+ )
+
+
+class DagBundleFileResponse(BaseModel):
+ """A file in a Dag bundle, as the Dag processor last saw it."""
+
+ relative_fileloc: str
+ dag_count: int = Field(description="Number of live Dags the file defines
that the caller may read.")
+ last_parsed_time: datetime | None = Field(
+ description="When the file was last parsed, or null if it has never
parsed successfully."
+ )
+ last_parse_duration: float | None = Field(
+ description="How long the last successful parse of the file took, in
seconds."
+ )
+ import_error_count: int | None = Field(
+ description=(
+ "Number of import errors recorded against the file, which is at
most one. Null when the "
+ "caller may not read import errors -- deliberately not zero, which
would read as a "
+ "file with nothing wrong."
+ )
+ )
+
+
+class DagBundleFileCollectionResponse(BaseModel):
+ """Dag bundle file collection response."""
+
+ dag_bundle_files: Iterable[DagBundleFileResponse]
+ total_entries: int
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 a9e9be01f4a..a8799ec436a 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
@@ -3417,6 +3417,152 @@ paths:
application/json:
schema:
$ref: '#/components/schemas/HTTPValidationError'
+ /api/v2/dagBundles/{bundle_name}:
+ get:
+ tags:
+ - Dag Bundle
+ summary: Get Dag Bundle
+ description: Get a Dag bundle.
+ operationId: get_dag_bundle
+ security:
+ - OAuth2PasswordBearer: []
+ - HTTPBearer: []
+ parameters:
+ - name: bundle_name
+ in: path
+ required: true
+ schema:
+ type: string
+ title: Bundle Name
+ responses:
+ '200':
+ description: Successful Response
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/DagBundleDetailResponse'
+ '401':
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ description: Unauthorized
+ '403':
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ description: Forbidden
+ '404':
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ description: Not Found
+ '422':
+ description: Validation Error
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPValidationError'
+ /api/v2/dagBundles/{bundle_name}/files:
+ get:
+ tags:
+ - Dag Bundle
+ summary: Get Dag Bundle Files
+ description: 'List the files in a Dag bundle, ordered by path.
+
+
+ A file is listed when the caller can read at least one Dag it defines.
A file
+ that recorded an
+
+ import error without registering any Dag is listed too, on the same
admin-by-default
+ terms as
+
+ ``GET /importErrors`` -- a file that fails before defining a Dag has
no Dag
+ to authorize on,
+
+ and it is the case this page most needs to show.
+
+
+ ``dag_count`` counts live Dags only. Dag rows are never deleted: a
file that
+ fails to import
+
+ has every Dag in it marked stale, and so does a file dropped from the
bundle.
+ A file therefore
+
+ stays listed on the strength of a live Dag *or* an import error, so a
file
+ that has just broken
+
+ does not vanish from the page someone opened to find out why, while a
file
+ deleted long ago
+
+ drops out instead of lingering forever.
+
+
+ Parse times come from the Dags in the file rather than the file
itself, which
+ is the only
+
+ record Airflow keeps: a file whose every Dag was removed keeps no
parse time
+ of its own.'
+ operationId: get_dag_bundle_files
+ security:
+ - OAuth2PasswordBearer: []
+ - HTTPBearer: []
+ parameters:
+ - name: bundle_name
+ in: path
+ required: true
+ schema:
+ type: string
+ title: Bundle Name
+ - name: limit
+ in: query
+ required: false
+ schema:
+ type: integer
+ minimum: 0
+ default: 50
+ title: Limit
+ - name: offset
+ in: query
+ required: false
+ schema:
+ type: integer
+ minimum: 0
+ default: 0
+ title: Offset
+ responses:
+ '200':
+ description: Successful Response
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/DagBundleFileCollectionResponse'
+ '401':
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ description: Unauthorized
+ '403':
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ description: Forbidden
+ '404':
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPExceptionResponse'
+ description: Not Found
+ '422':
+ description: Validation Error
+ content:
+ application/json:
+ schema:
+ $ref: '#/components/schemas/HTTPValidationError'
/api/v2/dagStats:
get:
tags:
@@ -14069,6 +14215,128 @@ components:
- total_entries
title: DagBundleCollectionResponse
description: Dag bundle collection response.
+ DagBundleDetailResponse:
+ properties:
+ name:
+ type: string
+ title: Name
+ active:
+ anyOf:
+ - type: boolean
+ - type: 'null'
+ title: Active
+ description: Whether the bundle is still present in this
deployment's configuration.
+ version:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Version
+ description: The latest version Airflow has seen for the bundle.
Null when
+ the bundle does not support versioning, or when no Dag processor
has refreshed
+ it successfully yet.
+ last_refreshed:
+ anyOf:
+ - type: string
+ format: date-time
+ - type: 'null'
+ title: Last Refreshed
+ description: When a Dag processor last successfully refreshed the
bundle.
+ It advances even when the version did not change, and a failed
refresh
+ leaves it untouched.
+ bundle_url:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Bundle Url
+ description: A link to view the bundle at ``version``, when one is
configured
+ and the caller may read Dag versions.
+ team_name:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Team Name
+ description: The team owning the bundle, in a multi-team deployment.
+ import_error_count:
+ anyOf:
+ - type: integer
+ - type: 'null'
+ title: Import Error Count
+ description: Number of Dag import errors recorded against this
bundle that
+ the caller is permitted to see, counted on the same terms as ``GET
/importErrors``.
+ Null when the caller may not read import errors.
+ dag_count:
+ type: integer
+ title: Dag Count
+ description: Number of live Dags recorded against the bundle that
the caller
+ is permitted to see, counted on the same terms as ``GET /dags``.
+ type: object
+ required:
+ - name
+ - active
+ - version
+ - last_refreshed
+ - bundle_url
+ - team_name
+ - import_error_count
+ - dag_count
+ title: DagBundleDetailResponse
+ description: Dag bundle serializer for the single-bundle response.
+ DagBundleFileCollectionResponse:
+ properties:
+ dag_bundle_files:
+ items:
+ $ref: '#/components/schemas/DagBundleFileResponse'
+ type: array
+ title: Dag Bundle Files
+ total_entries:
+ type: integer
+ title: Total Entries
+ type: object
+ required:
+ - dag_bundle_files
+ - total_entries
+ title: DagBundleFileCollectionResponse
+ description: Dag bundle file collection response.
+ DagBundleFileResponse:
+ properties:
+ relative_fileloc:
+ type: string
+ title: Relative Fileloc
+ dag_count:
+ type: integer
+ title: Dag Count
+ description: Number of live Dags the file defines that the caller
may read.
+ last_parsed_time:
+ anyOf:
+ - type: string
+ format: date-time
+ - type: 'null'
+ title: Last Parsed Time
+ description: When the file was last parsed, or null if it has never
parsed
+ successfully.
+ last_parse_duration:
+ anyOf:
+ - type: number
+ - type: 'null'
+ title: Last Parse Duration
+ description: How long the last successful parse of the file took, in
seconds.
+ import_error_count:
+ anyOf:
+ - type: integer
+ - type: 'null'
+ title: Import Error Count
+ description: Number of import errors recorded against the file,
which is
+ at most one. Null when the caller may not read import errors --
deliberately
+ not zero, which would read as a file with nothing wrong.
+ type: object
+ required:
+ - relative_fileloc
+ - dag_count
+ - last_parsed_time
+ - last_parse_duration
+ - import_error_count
+ title: DagBundleFileResponse
+ description: A file in a Dag bundle, as the Dag processor last saw it.
DagBundleResponse:
properties:
name:
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_bundles.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_bundles.py
index 296efd38ac9..a36ae59ea63 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_bundles.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_bundles.py
@@ -16,10 +16,11 @@
# under the License.
from __future__ import annotations
-from typing import TYPE_CHECKING, Annotated
+from operator import itemgetter
+from typing import TYPE_CHECKING, Annotated, Any, TypeVar
-from fastapi import Depends
-from sqlalchemy import and_, func, select
+from fastapi import Depends, HTTPException, status
+from sqlalchemy import and_, case, false, func, select
from airflow.api_fastapi.auth.managers.models.resource_details import
AccessView, DagAccessEntity
from airflow.api_fastapi.common.db.common import SessionDep, paginated_select
@@ -27,11 +28,16 @@ from airflow.api_fastapi.common.parameters import
QueryLimit, QueryOffset, SortP
from airflow.api_fastapi.common.router import AirflowRouter
from airflow.api_fastapi.core_api.datamodels.dag_bundles import (
DagBundleCollectionResponse,
+ DagBundleDetailResponse,
+ DagBundleFileCollectionResponse,
+ DagBundleFileResponse,
DagBundleResponse,
)
+from airflow.api_fastapi.core_api.openapi.exceptions import
create_openapi_http_exception_doc
from airflow.api_fastapi.core_api.security import (
AuthManagerDep,
GetUserDep,
+ PermittedDagBundleFilter,
ReadableDagBundlesFilterDep,
requires_access_dag,
)
@@ -50,6 +56,8 @@ if TYPE_CHECKING:
dag_bundles_router = AirflowRouter(tags=["Dag Bundle"], prefix="/dagBundles")
+_BundleResponseT = TypeVar("_BundleResponseT", bound=DagBundleResponse)
+
def _import_error_counts(
*,
@@ -58,6 +66,7 @@ def _import_error_counts(
auth_manager: BaseAuthManager,
user: BaseUser,
session: Session,
+ team_names: dict[str, str | None] | None = None,
) -> dict[str, int] | None:
"""
Count the import errors per bundle that this caller is allowed to know
about.
@@ -125,9 +134,10 @@ def _import_error_counts(
.group_by(ParseImportError.bundle_name)
).all()
if unregistered:
- team_names = DagBundleModel.get_team_names(
- [bundle_name for bundle_name, _ in unregistered if bundle_name],
session=session
- )
+ if team_names is None:
+ team_names = DagBundleModel.get_team_names(
+ [bundle_name for bundle_name, _ in unregistered if
bundle_name], session=session
+ )
for bundle_name, count in unregistered:
if bundle_name is None:
continue
@@ -141,6 +151,50 @@ def _import_error_counts(
return counts
+def _resolve_readable_bundle(
+ bundle_name: str,
+ readable_dag_bundles_filter: PermittedDagBundleFilter,
+ session: Session,
+) -> DagBundleModel:
+ """Load a bundle the caller is allowed to see, or raise 404 if there is no
such bundle."""
+ # 404 rather than 403 for a bundle that exists but holds no readable Dag,
so that the response
+ # does not confirm which bundle names a deployment has.
+ bundle = session.scalar(
+
readable_dag_bundles_filter.to_orm(select(DagBundleModel).where(DagBundleModel.name
== bundle_name))
+ )
+ if bundle is None:
+ raise HTTPException(status.HTTP_404_NOT_FOUND, f"Dag bundle with name
`{bundle_name}` was not found")
+ return bundle
+
+
+def _serialize_bundle(
+ response_class: type[_BundleResponseT],
+ bundle: DagBundleModel,
+ *,
+ team_name: str | None,
+ import_error_count: int | None,
+ show_bundle_url: bool,
+ **extra: Any,
+) -> _BundleResponseT:
+ """Build a bundle response, so the column-to-field mapping lives in one
place."""
+ return response_class(
+ name=bundle.name,
+ active=bundle.active,
+ version=bundle.version,
+ last_refreshed=bundle.last_refreshed,
+ bundle_url=bundle.render_url(bundle.version) if show_bundle_url else
None,
+ team_name=team_name,
+ import_error_count=import_error_count,
+ **extra,
+ )
+
+
+def _may_show_bundle_url(auth_manager: BaseAuthManager, user: BaseUser) ->
bool:
+ # A rendered bundle url is otherwise only reachable through ``GET
/dags/{dag_id}/dagVersions``,
+ # which additionally requires Dag *version* read, so withhold it from a
role without that.
+ return auth_manager.is_authorized_dag(method="GET",
access_entity=DagAccessEntity.VERSION, user=user)
+
+
@dag_bundles_router.get("",
dependencies=[Depends(requires_access_dag(method="GET"))])
def get_dag_bundles(
limit: QueryLimit,
@@ -192,26 +246,192 @@ def get_dag_bundles(
session=session,
)
- # A rendered bundle url is otherwise only reachable through ``GET
/dags/{dag_id}/dagVersions``,
- # which additionally requires Dag *version* read, so withhold it from a
role without that.
- show_bundle_url = auth_manager.is_authorized_dag(
- method="GET", access_entity=DagAccessEntity.VERSION, user=user
- )
+ show_bundle_url = _may_show_bundle_url(auth_manager, user)
return DagBundleCollectionResponse(
dag_bundles=[
- DagBundleResponse(
- name=bundle.name,
- active=bundle.active,
- version=bundle.version,
- last_refreshed=bundle.last_refreshed,
- bundle_url=bundle.render_url(bundle.version) if
show_bundle_url else None,
+ _serialize_bundle(
+ DagBundleResponse,
+ bundle,
team_name=team_names.get(bundle.name),
import_error_count=(
None if import_error_counts is None else
import_error_counts.get(bundle.name, 0)
),
+ show_bundle_url=show_bundle_url,
)
for bundle in bundles
],
total_entries=total_entries,
)
+
+
+@dag_bundles_router.get(
+ "/{bundle_name}",
+ responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
+ dependencies=[Depends(requires_access_dag(method="GET"))],
+)
+def get_dag_bundle(
+ bundle_name: str,
+ readable_dag_bundles_filter: ReadableDagBundlesFilterDep,
+ auth_manager: AuthManagerDep,
+ session: SessionDep,
+ user: GetUserDep,
+) -> DagBundleDetailResponse:
+ """Get a Dag bundle."""
+ bundle = _resolve_readable_bundle(bundle_name,
readable_dag_bundles_filter, session)
+ readable_dag_ids = readable_dag_bundles_filter.value or set()
+
+ # Stale rows are excluded, matching ``GET /dags``; see
``get_dag_bundle_files`` for why.
+ dag_count = session.scalar(
+ select(func.count(DagModel.dag_id)).where(
+ DagModel.bundle_name == bundle_name,
+ DagModel.dag_id.in_(readable_dag_ids),
+ DagModel.is_stale == false(),
+ )
+ )
+ # Resolved once and threaded through: ``_import_error_counts`` needs the
team to authorize
+ # the admin-gated view, and the response needs it again.
+ team_names = (
+ DagBundleModel.get_team_names([bundle_name], session=session)
+ if conf.getboolean("core", "multi_team")
+ else {}
+ )
+ import_error_counts = _import_error_counts(
+ bundle_names=[bundle_name],
+ readable_dag_ids=readable_dag_ids,
+ auth_manager=auth_manager,
+ user=user,
+ session=session,
+ team_names=team_names,
+ )
+
+ return _serialize_bundle(
+ DagBundleDetailResponse,
+ bundle,
+ team_name=team_names.get(bundle_name),
+ import_error_count=(None if import_error_counts is None else
import_error_counts.get(bundle_name, 0)),
+ show_bundle_url=_may_show_bundle_url(auth_manager, user),
+ dag_count=dag_count or 0,
+ )
+
+
+@dag_bundles_router.get(
+ "/{bundle_name}/files",
+ responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
+ dependencies=[Depends(requires_access_dag(method="GET"))],
+)
+def get_dag_bundle_files(
+ bundle_name: str,
+ limit: QueryLimit,
+ offset: QueryOffset,
+ readable_dag_bundles_filter: ReadableDagBundlesFilterDep,
+ auth_manager: AuthManagerDep,
+ session: SessionDep,
+ user: GetUserDep,
+) -> DagBundleFileCollectionResponse:
+ """
+ List the files in a Dag bundle, ordered by path.
+
+ A file is listed when the caller can read at least one Dag it defines. A
file that recorded an
+ import error without registering any Dag is listed too, on the same
admin-by-default terms as
+ ``GET /importErrors`` -- a file that fails before defining a Dag has no
Dag to authorize on,
+ and it is the case this page most needs to show.
+
+ ``dag_count`` counts live Dags only. Dag rows are never deleted: a file
that fails to import
+ has every Dag in it marked stale, and so does a file dropped from the
bundle. A file therefore
+ stays listed on the strength of a live Dag *or* an import error, so a file
that has just broken
+ does not vanish from the page someone opened to find out why, while a file
deleted long ago
+ drops out instead of lingering forever.
+
+ Parse times come from the Dags in the file rather than the file itself,
which is the only
+ record Airflow keeps: a file whose every Dag was removed keeps no parse
time of its own.
+ """
+ _resolve_readable_bundle(bundle_name, readable_dag_bundles_filter, session)
+ readable_dag_ids = readable_dag_bundles_filter.value or set()
+
+ # The duration is aggregated over live rows only. A Dag removed from a
file keeps its row,
+ # frozen at an older parse, so an unrestricted ``MAX`` would pair that
stale duration with the
+ # newer timestamp and report a parse that never happened.
+ files: dict[str, dict] = {
+ fileloc: {
+ "relative_fileloc": fileloc,
+ "dag_count": live_dag_count,
+ "last_parsed_time": last_parsed_time,
+ "last_parse_duration": last_parse_duration,
+ # Zero would read as "nothing wrong with this file"; the caller's
permission is
+ # decided below, and until then the count is unknown rather than
clean.
+ "import_error_count": None,
+ }
+ for fileloc, live_dag_count, last_parsed_time, last_parse_duration in
session.execute(
+ select(
+ DagModel.relative_fileloc,
+ func.count(case((DagModel.is_stale == false(),
DagModel.dag_id))),
+ func.max(DagModel.last_parsed_time),
+ func.max(case((DagModel.is_stale == false(),
DagModel.last_parse_duration))),
+ )
+ .where(
+ DagModel.bundle_name == bundle_name,
+ DagModel.dag_id.in_(readable_dag_ids),
+ DagModel.relative_fileloc.is_not(None),
+ )
+ .group_by(DagModel.relative_fileloc)
+ )
+ }
+
+ error_filenames: set[str] = set()
+ if auth_manager.authorize_view(access_view=AccessView.IMPORT_ERRORS,
user=user):
+ for row in files.values():
+ row["import_error_count"] = 0
+ error_filenames = {
+ filename
+ for filename in session.scalars(
+
select(ParseImportError.filename).where(ParseImportError.bundle_name ==
bundle_name)
+ )
+ if filename is not None
+ }
+ for filename in error_filenames & files.keys():
+ # ``import_error`` holds one row per file, so this is 0 or 1.
+ files[filename]["import_error_count"] = 1
+
+ unregistered = error_filenames - files.keys()
+ if unregistered and auth_manager.authorize_view(
+ access_view=AccessView.IMPORT_ERRORS_ALL,
+ user=user,
+ team_name=(
+ DagBundleModel.get_team_names([bundle_name],
session=session).get(bundle_name)
+ if conf.getboolean("core", "multi_team")
+ else None
+ ),
+ ):
+ # Queried only behind the permission that consumes it. A file with
no Dag row at all
+ # has no Dag to authorize on, so it needs the admin-by-default
view; a file whose Dags
+ # exist but are unreadable is excluded by being in this set.
+ files_with_any_dags = set(
+ session.scalars(
+
select(DagModel.relative_fileloc).where(DagModel.bundle_name ==
bundle_name).distinct()
+ )
+ )
+ for filename in unregistered - files_with_any_dags:
+ files[filename] = {
+ "relative_fileloc": filename,
+ "dag_count": 0,
+ "last_parsed_time": None,
+ "last_parse_duration": None,
+ "import_error_count": 1,
+ }
+
+ # A file with no live Dag and no error the caller can see is one the
bundle no longer has.
+ ordered = sorted(
+ (row for path, row in files.items() if row["dag_count"] > 0 or path in
error_filenames),
+ key=itemgetter("relative_fileloc"),
+ )
+ start = offset.value or 0
+ # ``limit`` is a non-negative int, so zero is a legal value meaning "no
rows" -- as
+ # ``Select.limit(0)`` gives every other endpoint. Testing it for
truthiness would read it as
+ # "unlimited" and return the whole list.
+ end = None if limit.value is None else start + limit.value
+
+ return DagBundleFileCollectionResponse(
+ dag_bundle_files=[DagBundleFileResponse(**row) for row in
ordered[start:end]],
+ total_entries=len(ordered),
+ )
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 0626c401617..05479e1f51d 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts
@@ -298,6 +298,20 @@ export const UseDagBundleServiceGetDagBundlesKeyFn = ({
limit, offset, orderBy }
offset?: number;
orderBy?: string[];
} = {}, queryKey?: Array<unknown>) => [useDagBundleServiceGetDagBundlesKey,
...(queryKey ?? [{ limit, offset, orderBy }])];
+export type DagBundleServiceGetDagBundleDefaultResponse =
Awaited<ReturnType<typeof DagBundleService.getDagBundle>>;
+export type DagBundleServiceGetDagBundleQueryResult<TData =
DagBundleServiceGetDagBundleDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
+export const useDagBundleServiceGetDagBundleKey =
"DagBundleServiceGetDagBundle";
+export const UseDagBundleServiceGetDagBundleKeyFn = ({ bundleName }: {
+ bundleName: string;
+}, queryKey?: Array<unknown>) => [useDagBundleServiceGetDagBundleKey,
...(queryKey ?? [{ bundleName }])];
+export type DagBundleServiceGetDagBundleFilesDefaultResponse =
Awaited<ReturnType<typeof DagBundleService.getDagBundleFiles>>;
+export type DagBundleServiceGetDagBundleFilesQueryResult<TData =
DagBundleServiceGetDagBundleFilesDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
+export const useDagBundleServiceGetDagBundleFilesKey =
"DagBundleServiceGetDagBundleFiles";
+export const UseDagBundleServiceGetDagBundleFilesKeyFn = ({ bundleName, limit,
offset }: {
+ bundleName: string;
+ limit?: number;
+ offset?: number;
+}, queryKey?: Array<unknown>) => [useDagBundleServiceGetDagBundleFilesKey,
...(queryKey ?? [{ bundleName, limit, offset }])];
export type DagStatsServiceGetDagStatsDefaultResponse =
Awaited<ReturnType<typeof DagStatsService.getDagStats>>;
export type DagStatsServiceGetDagStatsQueryResult<TData =
DagStatsServiceGetDagStatsDefaultResponse, TError = unknown> =
UseQueryResult<TData, TError>;
export const useDagStatsServiceGetDagStatsKey = "DagStatsServiceGetDagStats";
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 82c2cc61014..e1c55f73cf5 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts
@@ -585,6 +585,46 @@ export const ensureUseDagBundleServiceGetDagBundlesData =
(queryClient: QueryCli
orderBy?: string[];
} = {}) => queryClient.ensureQueryData({ queryKey:
Common.UseDagBundleServiceGetDagBundlesKeyFn({ limit, offset, orderBy }),
queryFn: () => DagBundleService.getDagBundles({ limit, offset, orderBy }) });
/**
+* Get Dag Bundle
+* Get a Dag bundle.
+* @param data The data for the request.
+* @param data.bundleName
+* @returns DagBundleDetailResponse Successful Response
+* @throws ApiError
+*/
+export const ensureUseDagBundleServiceGetDagBundleData = (queryClient:
QueryClient, { bundleName }: {
+ bundleName: string;
+}) => queryClient.ensureQueryData({ queryKey:
Common.UseDagBundleServiceGetDagBundleKeyFn({ bundleName }), queryFn: () =>
DagBundleService.getDagBundle({ bundleName }) });
+/**
+* Get Dag Bundle Files
+* List the files in a Dag bundle, ordered by path.
+*
+* A file is listed when the caller can read at least one Dag it defines. A
file that recorded an
+* import error without registering any Dag is listed too, on the same
admin-by-default terms as
+* ``GET /importErrors`` -- a file that fails before defining a Dag has no Dag
to authorize on,
+* and it is the case this page most needs to show.
+*
+* ``dag_count`` counts live Dags only. Dag rows are never deleted: a file that
fails to import
+* has every Dag in it marked stale, and so does a file dropped from the
bundle. A file therefore
+* stays listed on the strength of a live Dag *or* an import error, so a file
that has just broken
+* does not vanish from the page someone opened to find out why, while a file
deleted long ago
+* drops out instead of lingering forever.
+*
+* Parse times come from the Dags in the file rather than the file itself,
which is the only
+* record Airflow keeps: a file whose every Dag was removed keeps no parse time
of its own.
+* @param data The data for the request.
+* @param data.bundleName
+* @param data.limit
+* @param data.offset
+* @returns DagBundleFileCollectionResponse Successful Response
+* @throws ApiError
+*/
+export const ensureUseDagBundleServiceGetDagBundleFilesData = (queryClient:
QueryClient, { bundleName, limit, offset }: {
+ bundleName: string;
+ limit?: number;
+ offset?: number;
+}) => queryClient.ensureQueryData({ queryKey:
Common.UseDagBundleServiceGetDagBundleFilesKeyFn({ bundleName, limit, offset
}), queryFn: () => DagBundleService.getDagBundleFiles({ bundleName, limit,
offset }) });
+/**
* Get Dag Stats
* Get Dag statistics.
* @param data The data for the request.
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 56a373bd2c1..613336b861b 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts
@@ -585,6 +585,46 @@ export const prefetchUseDagBundleServiceGetDagBundles =
(queryClient: QueryClien
orderBy?: string[];
} = {}) => queryClient.prefetchQuery({ queryKey:
Common.UseDagBundleServiceGetDagBundlesKeyFn({ limit, offset, orderBy }),
queryFn: () => DagBundleService.getDagBundles({ limit, offset, orderBy }) });
/**
+* Get Dag Bundle
+* Get a Dag bundle.
+* @param data The data for the request.
+* @param data.bundleName
+* @returns DagBundleDetailResponse Successful Response
+* @throws ApiError
+*/
+export const prefetchUseDagBundleServiceGetDagBundle = (queryClient:
QueryClient, { bundleName }: {
+ bundleName: string;
+}) => queryClient.prefetchQuery({ queryKey:
Common.UseDagBundleServiceGetDagBundleKeyFn({ bundleName }), queryFn: () =>
DagBundleService.getDagBundle({ bundleName }) });
+/**
+* Get Dag Bundle Files
+* List the files in a Dag bundle, ordered by path.
+*
+* A file is listed when the caller can read at least one Dag it defines. A
file that recorded an
+* import error without registering any Dag is listed too, on the same
admin-by-default terms as
+* ``GET /importErrors`` -- a file that fails before defining a Dag has no Dag
to authorize on,
+* and it is the case this page most needs to show.
+*
+* ``dag_count`` counts live Dags only. Dag rows are never deleted: a file that
fails to import
+* has every Dag in it marked stale, and so does a file dropped from the
bundle. A file therefore
+* stays listed on the strength of a live Dag *or* an import error, so a file
that has just broken
+* does not vanish from the page someone opened to find out why, while a file
deleted long ago
+* drops out instead of lingering forever.
+*
+* Parse times come from the Dags in the file rather than the file itself,
which is the only
+* record Airflow keeps: a file whose every Dag was removed keeps no parse time
of its own.
+* @param data The data for the request.
+* @param data.bundleName
+* @param data.limit
+* @param data.offset
+* @returns DagBundleFileCollectionResponse Successful Response
+* @throws ApiError
+*/
+export const prefetchUseDagBundleServiceGetDagBundleFiles = (queryClient:
QueryClient, { bundleName, limit, offset }: {
+ bundleName: string;
+ limit?: number;
+ offset?: number;
+}) => queryClient.prefetchQuery({ queryKey:
Common.UseDagBundleServiceGetDagBundleFilesKeyFn({ bundleName, limit, offset
}), queryFn: () => DagBundleService.getDagBundleFiles({ bundleName, limit,
offset }) });
+/**
* Get Dag Stats
* Get Dag statistics.
* @param data The data for the request.
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 e92bdcccf7f..a20b8210436 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -585,6 +585,46 @@ export const useDagBundleServiceGetDagBundles = <TData =
Common.DagBundleService
orderBy?: string[];
} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseDagBundleServiceGetDagBundlesKeyFn({ limit, offset, orderBy },
queryKey), queryFn: () => DagBundleService.getDagBundles({ limit, offset,
orderBy }) as TData, ...options });
/**
+* Get Dag Bundle
+* Get a Dag bundle.
+* @param data The data for the request.
+* @param data.bundleName
+* @returns DagBundleDetailResponse Successful Response
+* @throws ApiError
+*/
+export const useDagBundleServiceGetDagBundle = <TData =
Common.DagBundleServiceGetDagBundleDefaultResponse, TError = unknown, TQueryKey
extends Array<unknown> = unknown[]>({ bundleName }: {
+ bundleName: string;
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseDagBundleServiceGetDagBundleKeyFn({ bundleName }, queryKey), queryFn:
() => DagBundleService.getDagBundle({ bundleName }) as TData, ...options });
+/**
+* Get Dag Bundle Files
+* List the files in a Dag bundle, ordered by path.
+*
+* A file is listed when the caller can read at least one Dag it defines. A
file that recorded an
+* import error without registering any Dag is listed too, on the same
admin-by-default terms as
+* ``GET /importErrors`` -- a file that fails before defining a Dag has no Dag
to authorize on,
+* and it is the case this page most needs to show.
+*
+* ``dag_count`` counts live Dags only. Dag rows are never deleted: a file that
fails to import
+* has every Dag in it marked stale, and so does a file dropped from the
bundle. A file therefore
+* stays listed on the strength of a live Dag *or* an import error, so a file
that has just broken
+* does not vanish from the page someone opened to find out why, while a file
deleted long ago
+* drops out instead of lingering forever.
+*
+* Parse times come from the Dags in the file rather than the file itself,
which is the only
+* record Airflow keeps: a file whose every Dag was removed keeps no parse time
of its own.
+* @param data The data for the request.
+* @param data.bundleName
+* @param data.limit
+* @param data.offset
+* @returns DagBundleFileCollectionResponse Successful Response
+* @throws ApiError
+*/
+export const useDagBundleServiceGetDagBundleFiles = <TData =
Common.DagBundleServiceGetDagBundleFilesDefaultResponse, TError = unknown,
TQueryKey extends Array<unknown> = unknown[]>({ bundleName, limit, offset }: {
+ bundleName: string;
+ limit?: number;
+ offset?: number;
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useQuery<TData, TError>({ queryKey:
Common.UseDagBundleServiceGetDagBundleFilesKeyFn({ bundleName, limit, offset },
queryKey), queryFn: () => DagBundleService.getDagBundleFiles({ bundleName,
limit, offset }) as TData, ...options });
+/**
* Get Dag Stats
* Get Dag statistics.
* @param data The data for the request.
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 8a8da0566e0..9add7eba93e 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts
@@ -585,6 +585,46 @@ export const useDagBundleServiceGetDagBundlesSuspense =
<TData = Common.DagBundl
orderBy?: string[];
} = {}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseDagBundleServiceGetDagBundlesKeyFn({ limit, offset, orderBy },
queryKey), queryFn: () => DagBundleService.getDagBundles({ limit, offset,
orderBy }) as TData, ...options });
/**
+* Get Dag Bundle
+* Get a Dag bundle.
+* @param data The data for the request.
+* @param data.bundleName
+* @returns DagBundleDetailResponse Successful Response
+* @throws ApiError
+*/
+export const useDagBundleServiceGetDagBundleSuspense = <TData =
Common.DagBundleServiceGetDagBundleDefaultResponse, TError = unknown, TQueryKey
extends Array<unknown> = unknown[]>({ bundleName }: {
+ bundleName: string;
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseDagBundleServiceGetDagBundleKeyFn({ bundleName }, queryKey), queryFn:
() => DagBundleService.getDagBundle({ bundleName }) as TData, ...options });
+/**
+* Get Dag Bundle Files
+* List the files in a Dag bundle, ordered by path.
+*
+* A file is listed when the caller can read at least one Dag it defines. A
file that recorded an
+* import error without registering any Dag is listed too, on the same
admin-by-default terms as
+* ``GET /importErrors`` -- a file that fails before defining a Dag has no Dag
to authorize on,
+* and it is the case this page most needs to show.
+*
+* ``dag_count`` counts live Dags only. Dag rows are never deleted: a file that
fails to import
+* has every Dag in it marked stale, and so does a file dropped from the
bundle. A file therefore
+* stays listed on the strength of a live Dag *or* an import error, so a file
that has just broken
+* does not vanish from the page someone opened to find out why, while a file
deleted long ago
+* drops out instead of lingering forever.
+*
+* Parse times come from the Dags in the file rather than the file itself,
which is the only
+* record Airflow keeps: a file whose every Dag was removed keeps no parse time
of its own.
+* @param data The data for the request.
+* @param data.bundleName
+* @param data.limit
+* @param data.offset
+* @returns DagBundleFileCollectionResponse Successful Response
+* @throws ApiError
+*/
+export const useDagBundleServiceGetDagBundleFilesSuspense = <TData =
Common.DagBundleServiceGetDagBundleFilesDefaultResponse, TError = unknown,
TQueryKey extends Array<unknown> = unknown[]>({ bundleName, limit, offset }: {
+ bundleName: string;
+ limit?: number;
+ offset?: number;
+}, queryKey?: TQueryKey, options?: Omit<UseQueryOptions<TData, TError>,
"queryKey" | "queryFn">) => useSuspenseQuery<TData, TError>({ queryKey:
Common.UseDagBundleServiceGetDagBundleFilesKeyFn({ bundleName, limit, offset },
queryKey), queryFn: () => DagBundleService.getDagBundleFiles({ bundleName,
limit, offset }) as TData, ...options });
+/**
* Get Dag Stats
* Get Dag statistics.
* @param data The data for the request.
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 566cd1dba51..2aca8e5aeab 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
@@ -4454,6 +4454,172 @@ export const $DagBundleCollectionResponse = {
description: 'Dag bundle collection response.'
} as const;
+export const $DagBundleDetailResponse = {
+ properties: {
+ name: {
+ type: 'string',
+ title: 'Name'
+ },
+ active: {
+ anyOf: [
+ {
+ type: 'boolean'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Active',
+ description: "Whether the bundle is still present in this
deployment's configuration."
+ },
+ version: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Version',
+ description: 'The latest version Airflow has seen for the bundle.
Null when the bundle does not support versioning, or when no Dag processor has
refreshed it successfully yet.'
+ },
+ last_refreshed: {
+ anyOf: [
+ {
+ type: 'string',
+ format: 'date-time'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Last Refreshed',
+ description: 'When a Dag processor last successfully refreshed the
bundle. It advances even when the version did not change, and a failed refresh
leaves it untouched.'
+ },
+ bundle_url: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Bundle Url',
+ description: 'A link to view the bundle at ``version``, when one
is configured and the caller may read Dag versions.'
+ },
+ team_name: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Team Name',
+ description: 'The team owning the bundle, in a multi-team
deployment.'
+ },
+ import_error_count: {
+ anyOf: [
+ {
+ type: 'integer'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Import Error Count',
+ description: 'Number of Dag import errors recorded against this
bundle that the caller is permitted to see, counted on the same terms as ``GET
/importErrors``. Null when the caller may not read import errors.'
+ },
+ dag_count: {
+ type: 'integer',
+ title: 'Dag Count',
+ description: 'Number of live Dags recorded against the bundle that
the caller is permitted to see, counted on the same terms as ``GET /dags``.'
+ }
+ },
+ type: 'object',
+ required: ['name', 'active', 'version', 'last_refreshed', 'bundle_url',
'team_name', 'import_error_count', 'dag_count'],
+ title: 'DagBundleDetailResponse',
+ description: 'Dag bundle serializer for the single-bundle response.'
+} as const;
+
+export const $DagBundleFileCollectionResponse = {
+ properties: {
+ dag_bundle_files: {
+ items: {
+ '$ref': '#/components/schemas/DagBundleFileResponse'
+ },
+ type: 'array',
+ title: 'Dag Bundle Files'
+ },
+ total_entries: {
+ type: 'integer',
+ title: 'Total Entries'
+ }
+ },
+ type: 'object',
+ required: ['dag_bundle_files', 'total_entries'],
+ title: 'DagBundleFileCollectionResponse',
+ description: 'Dag bundle file collection response.'
+} as const;
+
+export const $DagBundleFileResponse = {
+ properties: {
+ relative_fileloc: {
+ type: 'string',
+ title: 'Relative Fileloc'
+ },
+ dag_count: {
+ type: 'integer',
+ title: 'Dag Count',
+ description: 'Number of live Dags the file defines that the caller
may read.'
+ },
+ last_parsed_time: {
+ anyOf: [
+ {
+ type: 'string',
+ format: 'date-time'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Last Parsed Time',
+ description: 'When the file was last parsed, or null if it has
never parsed successfully.'
+ },
+ last_parse_duration: {
+ anyOf: [
+ {
+ type: 'number'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Last Parse Duration',
+ description: 'How long the last successful parse of the file took,
in seconds.'
+ },
+ import_error_count: {
+ anyOf: [
+ {
+ type: 'integer'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Import Error Count',
+ description: 'Number of import errors recorded against the file,
which is at most one. Null when the caller may not read import errors --
deliberately not zero, which would read as a file with nothing wrong.'
+ }
+ },
+ type: 'object',
+ required: ['relative_fileloc', 'dag_count', 'last_parsed_time',
'last_parse_duration', 'import_error_count'],
+ title: 'DagBundleFileResponse',
+ description: 'A file in a Dag bundle, as the Dag processor last saw it.'
+} as const;
+
export const $DagBundleResponse = {
properties: {
name: {
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 baf23138c35..c564e43f2e7 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
@@ -3,7 +3,7 @@
import type { CancelablePromise } from './core/CancelablePromise';
import { OpenAPI } from './core/OpenAPI';
import { request as __request } from './core/request';
-import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData,
GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse,
GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData,
CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse,
GetAssetQueuedEventsData, GetAssetQueuedEventsResponse,
DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData,
GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse,
Dele [...]
+import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData,
GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse,
GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData,
CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse,
GetAssetQueuedEventsData, GetAssetQueuedEventsResponse,
DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData,
GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse,
Dele [...]
export class AssetService {
/**
@@ -1598,6 +1598,74 @@ export class DagBundleService {
});
}
+ /**
+ * Get Dag Bundle
+ * Get a Dag bundle.
+ * @param data The data for the request.
+ * @param data.bundleName
+ * @returns DagBundleDetailResponse Successful Response
+ * @throws ApiError
+ */
+ public static getDagBundle(data: GetDagBundleData):
CancelablePromise<GetDagBundleResponse> {
+ return __request(OpenAPI, {
+ method: 'GET',
+ url: '/api/v2/dagBundles/{bundle_name}',
+ path: {
+ bundle_name: data.bundleName
+ },
+ errors: {
+ 401: 'Unauthorized',
+ 403: 'Forbidden',
+ 404: 'Not Found',
+ 422: 'Validation Error'
+ }
+ });
+ }
+
+ /**
+ * Get Dag Bundle Files
+ * List the files in a Dag bundle, ordered by path.
+ *
+ * A file is listed when the caller can read at least one Dag it defines.
A file that recorded an
+ * import error without registering any Dag is listed too, on the same
admin-by-default terms as
+ * ``GET /importErrors`` -- a file that fails before defining a Dag has no
Dag to authorize on,
+ * and it is the case this page most needs to show.
+ *
+ * ``dag_count`` counts live Dags only. Dag rows are never deleted: a file
that fails to import
+ * has every Dag in it marked stale, and so does a file dropped from the
bundle. A file therefore
+ * stays listed on the strength of a live Dag *or* an import error, so a
file that has just broken
+ * does not vanish from the page someone opened to find out why, while a
file deleted long ago
+ * drops out instead of lingering forever.
+ *
+ * Parse times come from the Dags in the file rather than the file itself,
which is the only
+ * record Airflow keeps: a file whose every Dag was removed keeps no parse
time of its own.
+ * @param data The data for the request.
+ * @param data.bundleName
+ * @param data.limit
+ * @param data.offset
+ * @returns DagBundleFileCollectionResponse Successful Response
+ * @throws ApiError
+ */
+ public static getDagBundleFiles(data: GetDagBundleFilesData):
CancelablePromise<GetDagBundleFilesResponse> {
+ return __request(OpenAPI, {
+ method: 'GET',
+ url: '/api/v2/dagBundles/{bundle_name}/files',
+ path: {
+ bundle_name: data.bundleName
+ },
+ query: {
+ limit: data.limit,
+ offset: data.offset
+ },
+ errors: {
+ 401: 'Unauthorized',
+ 403: 'Forbidden',
+ 404: 'Not Found',
+ 422: 'Validation Error'
+ }
+ });
+ }
+
}
export class DagStatsService {
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 91ddaf68f7d..8cf3d7bb7d5 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
@@ -1151,6 +1151,72 @@ export type DagBundleCollectionResponse = {
total_entries: number;
};
+/**
+ * Dag bundle serializer for the single-bundle response.
+ */
+export type DagBundleDetailResponse = {
+ name: string;
+ /**
+ * Whether the bundle is still present in this deployment's configuration.
+ */
+ active: boolean | null;
+ /**
+ * The latest version Airflow has seen for the bundle. Null when the
bundle does not support versioning, or when no Dag processor has refreshed it
successfully yet.
+ */
+ version: string | null;
+ /**
+ * When a Dag processor last successfully refreshed the bundle. It
advances even when the version did not change, and a failed refresh leaves it
untouched.
+ */
+ last_refreshed: string | null;
+ /**
+ * A link to view the bundle at ``version``, when one is configured and
the caller may read Dag versions.
+ */
+ bundle_url: string | null;
+ /**
+ * The team owning the bundle, in a multi-team deployment.
+ */
+ team_name: string | null;
+ /**
+ * Number of Dag import errors recorded against this bundle that the
caller is permitted to see, counted on the same terms as ``GET /importErrors``.
Null when the caller may not read import errors.
+ */
+ import_error_count: number | null;
+ /**
+ * Number of live Dags recorded against the bundle that the caller is
permitted to see, counted on the same terms as ``GET /dags``.
+ */
+ dag_count: number;
+};
+
+/**
+ * Dag bundle file collection response.
+ */
+export type DagBundleFileCollectionResponse = {
+ dag_bundle_files: Array<DagBundleFileResponse>;
+ total_entries: number;
+};
+
+/**
+ * A file in a Dag bundle, as the Dag processor last saw it.
+ */
+export type DagBundleFileResponse = {
+ relative_fileloc: string;
+ /**
+ * Number of live Dags the file defines that the caller may read.
+ */
+ dag_count: number;
+ /**
+ * When the file was last parsed, or null if it has never parsed
successfully.
+ */
+ last_parsed_time: string | null;
+ /**
+ * How long the last successful parse of the file took, in seconds.
+ */
+ last_parse_duration: number | null;
+ /**
+ * Number of import errors recorded against the file, which is at most
one. Null when the caller may not read import errors -- deliberately not zero,
which would read as a file with nothing wrong.
+ */
+ import_error_count: number | null;
+};
+
/**
* Dag bundle serializer for responses.
*/
@@ -3546,6 +3612,20 @@ export type GetDagBundlesData = {
export type GetDagBundlesResponse = DagBundleCollectionResponse;
+export type GetDagBundleData = {
+ bundleName: string;
+};
+
+export type GetDagBundleResponse = DagBundleDetailResponse;
+
+export type GetDagBundleFilesData = {
+ bundleName: string;
+ limit?: number;
+ offset?: number;
+};
+
+export type GetDagBundleFilesResponse = DagBundleFileCollectionResponse;
+
export type GetDagStatsData = {
dagIds?: Array<(string)>;
};
@@ -6312,6 +6392,60 @@ export type $OpenApiTs = {
};
};
};
+ '/api/v2/dagBundles/{bundle_name}': {
+ get: {
+ req: GetDagBundleData;
+ res: {
+ /**
+ * Successful Response
+ */
+ 200: DagBundleDetailResponse;
+ /**
+ * Unauthorized
+ */
+ 401: HTTPExceptionResponse;
+ /**
+ * Forbidden
+ */
+ 403: HTTPExceptionResponse;
+ /**
+ * Not Found
+ */
+ 404: HTTPExceptionResponse;
+ /**
+ * Validation Error
+ */
+ 422: HTTPValidationError;
+ };
+ };
+ };
+ '/api/v2/dagBundles/{bundle_name}/files': {
+ get: {
+ req: GetDagBundleFilesData;
+ res: {
+ /**
+ * Successful Response
+ */
+ 200: DagBundleFileCollectionResponse;
+ /**
+ * Unauthorized
+ */
+ 401: HTTPExceptionResponse;
+ /**
+ * Forbidden
+ */
+ 403: HTTPExceptionResponse;
+ /**
+ * Not Found
+ */
+ 404: HTTPExceptionResponse;
+ /**
+ * Validation Error
+ */
+ 422: HTTPValidationError;
+ };
+ };
+ };
'/api/v2/dagStats': {
get: {
req: GetDagStatsData;
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/browse.json
b/airflow-core/src/airflow/ui/public/i18n/locales/en/browse.json
index 4e8910d8a9c..48ea7aa8155 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/browse.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/browse.json
@@ -12,8 +12,6 @@
},
"dagBundles": {
"active": "Active",
- "bundle_one": "Dag Bundle",
- "bundle_other": "Dag Bundles",
"columns": {
"active": "Active",
"importErrors": "Import Errors",
@@ -21,6 +19,22 @@
"name": "Name",
"version": "Version"
},
+ "detail": {
+ "status": "Status"
+ },
+ "files": {
+ "columns": {
+ "file": "File",
+ "lastParsed": "Last Parsed",
+ "parseDuration": "Parse Duration"
+ },
+ "file_one": "File",
+ "file_other": "Files",
+ "importErrorGone": "This import error is no longer recorded. It was
probably fixed by a reparse.",
+ "neverParsed": "Never",
+ "showImportError_one": "View error",
+ "showImportError_other": "View {{count}} errors"
+ },
"inactive": "Inactive",
"neverRefreshed": "Never",
"notRefreshedYet": "Not refreshed yet",
diff --git a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
index 2c4c451cf2d..8bc87e77e4a 100644
--- a/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
+++ b/airflow-core/src/airflow/ui/public/i18n/locales/en/common.json
@@ -40,6 +40,8 @@
"createdAssetEvent_other": "Created $t(common:assetEvent_other)",
"dag_one": "Dag",
"dag_other": "Dags",
+ "dagBundle_one": "Dag Bundle",
+ "dagBundle_other": "Dag Bundles",
"dagDetails": {
"activeRuns": "Active Runs",
"catchup": "Catchup",
diff --git a/airflow-core/src/airflow/ui/src/components/DagBundleVersion.tsx
b/airflow-core/src/airflow/ui/src/components/DagBundleVersion.tsx
new file mode 100644
index 00000000000..9163b46e49d
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/components/DagBundleVersion.tsx
@@ -0,0 +1,63 @@
+/*!
+ * 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.
+ */
+import { Link, Text } from "@chakra-ui/react";
+import { useTranslation } from "react-i18next";
+
+import { Tooltip } from "src/system-components";
+
+// A git bundle stores the full 40-char hexsha, so show the prefix an author
recognises and keep
+// the whole value on hover.
+const SHORT_VERSION_LENGTH = 7;
+
+type Props = {
+ readonly bundleUrl: string | null;
+ readonly lastRefreshed: string | null;
+ readonly version: string | null;
+};
+
+/** The version a bundle currently holds, linked to the commit it names where
that is known. */
+export const DagBundleVersion = ({ bundleUrl, lastRefreshed, version }: Props)
=> {
+ const { t: translate } = useTranslation("browse");
+
+ if (version === null) {
+ // A versioning-capable bundle also reports null until its first
successful refresh, so an
+ // absent last_refreshed is what separates "not refreshed yet" from "not
versioned".
+ return (
+ <Text color="fg.muted">
+ {lastRefreshed === null
+ ? translate("dagBundles.notRefreshedYet")
+ : translate("dagBundles.notVersioned")}
+ </Text>
+ );
+ }
+
+ const short = version.slice(0, SHORT_VERSION_LENGTH);
+
+ return (
+ <Tooltip content={version}>
+ {bundleUrl === null ? (
+ <Text fontFamily="mono">{short}</Text>
+ ) : (
+ <Link color="fg.info" fontFamily="mono" href={bundleUrl}
rel="noreferrer" target="_blank">
+ {short}
+ </Link>
+ )}
+ </Tooltip>
+ );
+};
diff --git
a/airflow-core/src/airflow/ui/src/components/DagDeactivatedBanner.tsx
b/airflow-core/src/airflow/ui/src/components/DagDeactivatedBanner.tsx
index 86ade22c0ef..acaa52e973f 100644
--- a/airflow-core/src/airflow/ui/src/components/DagDeactivatedBanner.tsx
+++ b/airflow-core/src/airflow/ui/src/components/DagDeactivatedBanner.tsx
@@ -23,7 +23,7 @@ import { useParams } from "react-router-dom";
import { useDagServiceGetDag, useImportErrorServiceGetImportErrors } from
"openapi/queries";
-import { DagImportErrorModal } from "src/pages/Dag/DagImportErrorModal";
+import { DagImportErrorModal } from "./DagImportErrorModal";
export const DagDeactivatedBanner = () => {
const { t: translate } = useTranslation(["dag", "dashboard"]);
diff --git a/airflow-core/src/airflow/ui/src/pages/Dag/DagImportErrorModal.tsx
b/airflow-core/src/airflow/ui/src/components/DagImportErrorModal.tsx
similarity index 100%
rename from airflow-core/src/airflow/ui/src/pages/Dag/DagImportErrorModal.tsx
rename to airflow-core/src/airflow/ui/src/components/DagImportErrorModal.tsx
diff --git a/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx
b/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx
index 90e9af5126b..5a6dce6312a 100644
--- a/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx
+++ b/airflow-core/src/airflow/ui/src/components/HeaderCard.tsx
@@ -34,7 +34,7 @@ type Props = {
readonly stats: Array<{ key?: string; label: string; value: ReactNode |
string }>;
readonly subTitle?: ReactNode | string;
readonly title: ReactNode | string;
- readonly type: "asset" | "dag" | "dagRun" | "task" | "taskGroup" |
"taskInstance";
+ readonly type: "asset" | "dag" | "dagBundle" | "dagRun" | "task" |
"taskGroup" | "taskInstance";
};
export const HeaderCard = ({ actions, icon, state, stats, subTitle, title,
type }: Props) => {
diff --git a/airflow-core/src/airflow/ui/src/components/ImportErrorCount.tsx
b/airflow-core/src/airflow/ui/src/components/ImportErrorCount.tsx
new file mode 100644
index 00000000000..8459b9cf4d2
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/components/ImportErrorCount.tsx
@@ -0,0 +1,43 @@
+/*!
+ * 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.
+ */
+import { Badge, Text } from "@chakra-ui/react";
+
+type Props = {
+ readonly count: number | null;
+};
+
+/**
+ * An import-error count, flagged red only when there is something to flag.
+ *
+ * Null is not zero: it means the caller may not read import errors, so it
renders as unknown
+ * rather than as a clean bill of health.
+ */
+export const ImportErrorCount = ({ count }: Props) => {
+ if (count === null) {
+ return <Text color="fg.muted">-</Text>;
+ }
+
+ return count === 0 ? (
+ <Text color="fg.muted">0</Text>
+ ) : (
+ <Badge colorPalette="failed" variant="solid">
+ {count}
+ </Badge>
+ );
+};
diff --git a/airflow-core/src/airflow/ui/src/pages/DagBundle/BundleFiles.tsx
b/airflow-core/src/airflow/ui/src/pages/DagBundle/BundleFiles.tsx
new file mode 100644
index 00000000000..f8308d0571f
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/DagBundle/BundleFiles.tsx
@@ -0,0 +1,144 @@
+/*!
+ * 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.
+ */
+import { useState } from "react";
+
+import { Box, Button, Text } from "@chakra-ui/react";
+import type { ColumnDef } from "@tanstack/react-table";
+import type { TFunction } from "i18next";
+import { useTranslation } from "react-i18next";
+import { LuFileWarning } from "react-icons/lu";
+
+import { useDagBundleServiceGetDagBundleFiles } from "openapi/queries";
+import type { DagBundleFileResponse } from "openapi/requests/types.gen";
+
+import { DataTable } from "src/components/DataTable";
+import { useTableURLState } from "src/components/DataTable/useTableUrlState";
+import { DurationCell } from "src/components/DurationCell";
+import { ErrorAlert } from "src/components/ErrorAlert";
+import { ImportErrorCount } from "src/components/ImportErrorCount";
+import Time from "src/components/Time";
+
+import { useDagBundleRefetchInterval } from
"src/queries/useDagBundleRefetchInterval";
+
+import { FileImportError } from "./FileImportError";
+
+type FileRow = { row: { original: DagBundleFileResponse } };
+
+// The endpoint assembles rows from two sources and orders them by path, so it
takes no sort
+// parameter and the table offers none.
+const createColumns = (
+ translate: TFunction,
+ onShowImportError: (relativeFileloc: string) => void,
+): Array<ColumnDef<DagBundleFileResponse>> => [
+ {
+ accessorKey: "relative_fileloc",
+ cell: ({ row: { original } }: FileRow) => <Text
fontFamily="mono">{original.relative_fileloc}</Text>,
+ enableSorting: false,
+ header: translate("browse:dagBundles.files.columns.file"),
+ },
+ {
+ accessorKey: "dag_count",
+ enableSorting: false,
+ header: translate("common:dag_other"),
+ },
+ {
+ accessorKey: "last_parsed_time",
+ cell: ({ row: { original } }: FileRow) =>
+ original.last_parsed_time === null ? (
+ <Text
color="fg.muted">{translate("browse:dagBundles.files.neverParsed")}</Text>
+ ) : (
+ <Time datetime={original.last_parsed_time} />
+ ),
+ enableSorting: false,
+ header: translate("browse:dagBundles.files.columns.lastParsed"),
+ },
+ {
+ accessorKey: "last_parse_duration",
+ cell: ({ row: { original } }: FileRow) => <DurationCell
duration={original.last_parse_duration} />,
+ enableSorting: false,
+ header: translate("browse:dagBundles.files.columns.parseDuration"),
+ },
+ {
+ accessorKey: "import_error_count",
+ cell: ({ row: { original } }: FileRow) =>
+ // Only worth clicking when there is an error behind the count. A bare
badge reads as a
+ // status rather than a control, so the clickable form carries an icon
and its own label.
+ original.import_error_count === null || original.import_error_count ===
0 ? (
+ <ImportErrorCount count={original.import_error_count} />
+ ) : (
+ <Button
+ colorPalette="failed"
+ onClick={() => {
+ onShowImportError(original.relative_fileloc);
+ }}
+ size="xs"
+ variant="outline"
+ >
+ <LuFileWarning />
+ {translate("browse:dagBundles.files.showImportError", {
+ count: original.import_error_count,
+ })}
+ </Button>
+ ),
+ enableSorting: false,
+ header: translate("browse:dagBundles.columns.importErrors"),
+ },
+];
+
+export const BundleFiles = ({ bundleName }: { readonly bundleName: string })
=> {
+ const { t: translate } = useTranslation(["browse", "common"]);
+ const refetchInterval = useDagBundleRefetchInterval();
+
+ const { setTableURLState, tableURLState } = useTableURLState();
+ const { pagination } = tableURLState;
+ const [errorFileloc, setErrorFileloc] = useState<string |
undefined>(undefined);
+
+ const { data, error, isFetching, isLoading } =
useDagBundleServiceGetDagBundleFiles(
+ {
+ bundleName,
+ limit: pagination.pageSize,
+ offset: pagination.pageIndex * pagination.pageSize,
+ },
+ undefined,
+ { refetchInterval, refetchOnWindowFocus: refetchInterval !== false },
+ );
+
+ return (
+ <Box pt={2}>
+ <FileImportError
+ bundleName={bundleName}
+ onClose={() => {
+ setErrorFileloc(undefined);
+ }}
+ relativeFileloc={errorFileloc}
+ />
+ <DataTable
+ columns={createColumns(translate, setErrorFileloc)}
+ data={data?.dag_bundle_files ?? []}
+ errorMessage={<ErrorAlert error={error} />}
+ initialState={tableURLState}
+ isFetching={isFetching}
+ isLoading={isLoading}
+ modelName="browse:dagBundles.files.file"
+ onStateChange={setTableURLState}
+ total={data?.total_entries}
+ />
+ </Box>
+ );
+};
diff --git a/airflow-core/src/airflow/ui/src/pages/DagBundle/DagBundle.tsx
b/airflow-core/src/airflow/ui/src/pages/DagBundle/DagBundle.tsx
new file mode 100644
index 00000000000..94d88de779c
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/DagBundle/DagBundle.tsx
@@ -0,0 +1,62 @@
+/*!
+ * 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.
+ */
+import { Box } from "@chakra-ui/react";
+import { useParams } from "react-router-dom";
+
+import { useDagBundleServiceGetDagBundle } from "openapi/queries";
+
+import { ProgressBar } from "src/system-components";
+
+import { ErrorAlert } from "src/components/ErrorAlert";
+
+import { useDagBundleRefetchInterval } from
"src/queries/useDagBundleRefetchInterval";
+import { useDocumentTitle } from "src/utils";
+
+import { BundleFiles } from "./BundleFiles";
+import { Header } from "./Header";
+
+export const DagBundle = () => {
+ // The route segment is required, so react-router cannot match this without
a name.
+ const { bundleName = "" } = useParams();
+ const refetchInterval = useDagBundleRefetchInterval();
+
+ const { data, error, isLoading } = useDagBundleServiceGetDagBundle({
bundleName }, undefined, {
+ refetchInterval,
+ // "Has my deploy landed yet" is exactly what a returning tab is asking,
and the version and
+ // last-refreshed fields that answer it live on this query.
+ refetchOnWindowFocus: refetchInterval !== false,
+ });
+
+ useDocumentTitle(bundleName);
+
+ return (
+ <Box p={2}>
+ <ErrorAlert error={error} />
+ <ProgressBar size="xs" visibility={isLoading ? "visible" : "hidden"} />
+ {/* Nothing below renders without a bundle: a header built from
undefined asserts an
+ inactive, never-refreshed bundle, and the file list would repeat the
same 404. */}
+ {data === undefined ? undefined : (
+ <>
+ <Header bundle={data} />
+ <BundleFiles bundleName={bundleName} />
+ </>
+ )}
+ </Box>
+ );
+};
diff --git
a/airflow-core/src/airflow/ui/src/pages/DagBundle/FileImportError.tsx
b/airflow-core/src/airflow/ui/src/pages/DagBundle/FileImportError.tsx
new file mode 100644
index 00000000000..ab3129be152
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/DagBundle/FileImportError.tsx
@@ -0,0 +1,78 @@
+/*!
+ * 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.
+ */
+import { Heading, HStack, Text } from "@chakra-ui/react";
+import { useTranslation } from "react-i18next";
+import { LuFileWarning } from "react-icons/lu";
+
+import { useImportErrorServiceGetImportErrors } from "openapi/queries";
+
+import { Modal } from "src/system-components";
+
+import { DagImportErrorModal } from "src/components/DagImportErrorModal";
+
+type Props = {
+ readonly bundleName: string;
+ readonly onClose: () => void;
+ readonly relativeFileloc: string | undefined;
+};
+
+/**
+ * Fetches the import error behind a file's count and hands it to the shared
modal.
+ *
+ * Read through ``GET /importErrors`` rather than returned with the file list:
a stack trace is
+ * unbounded and the file list is polled, and that endpoint already decides
who may read an import
+ * error. A dialog rather than a page of its own, because a single stack trace
is not worth an
+ * addressable route.
+ */
+export const FileImportError = ({ bundleName, onClose, relativeFileloc }:
Props) => {
+ const { t: translate } = useTranslation("browse");
+
+ const { data, isLoading } = useImportErrorServiceGetImportErrors(
+ { bundleName, filename: relativeFileloc ?? "" },
+ undefined,
+ { enabled: relativeFileloc !== undefined },
+ );
+
+ const [importError] = data?.import_errors ?? [];
+
+ if (importError !== undefined) {
+ return (
+ <DagImportErrorModal importError={importError} onClose={onClose}
open={relativeFileloc !== undefined} />
+ );
+ }
+
+ // The table polls, so the error can be gone by the time its badge is
clicked. Say so rather
+ // than leaving the click with no visible effect.
+ return (
+ <Modal
+ headerProps={{
+ children: (
+ <HStack gap={2}>
+ <LuFileWarning />
+ <Heading fontSize="lg">{relativeFileloc}</Heading>
+ </HStack>
+ ),
+ }}
+ onOpenChange={onClose}
+ open={relativeFileloc !== undefined && !isLoading}
+ >
+ <Text
color="fg.muted">{translate("dagBundles.files.importErrorGone")}</Text>
+ </Modal>
+ );
+};
diff --git a/airflow-core/src/airflow/ui/src/pages/DagBundle/Header.tsx
b/airflow-core/src/airflow/ui/src/pages/DagBundle/Header.tsx
new file mode 100644
index 00000000000..f16c588e841
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/DagBundle/Header.tsx
@@ -0,0 +1,84 @@
+/*!
+ * 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.
+ */
+import { Badge } from "@chakra-ui/react";
+import { useTranslation } from "react-i18next";
+import { FiPackage } from "react-icons/fi";
+
+import type { DagBundleDetailResponse } from "openapi/requests/types.gen";
+
+import { DagBundleVersion } from "src/components/DagBundleVersion";
+import { HeaderCard } from "src/components/HeaderCard";
+import { ImportErrorCount } from "src/components/ImportErrorCount";
+import { TeamName } from "src/components/TeamName";
+import Time from "src/components/Time";
+
+import { useShowTeam } from "src/hooks/useShowTeam";
+
+export const Header = ({ bundle }: { readonly bundle: DagBundleDetailResponse
}) => {
+ const { t: translate } = useTranslation(["browse", "common"]);
+ // Teams exist only in multi-team deployments, and a bundle need not belong
to one; an
+ // empty "Team" stat reads as a broken page.
+ const showTeam = useShowTeam(bundle.team_name);
+
+ const stats = [
+ {
+ label: translate("browse:dagBundles.detail.status"),
+ value:
+ bundle.active === true ? (
+ <Badge
colorPalette="success">{translate("browse:dagBundles.active")}</Badge>
+ ) : (
+ <Badge
colorPalette="gray">{translate("browse:dagBundles.inactive")}</Badge>
+ ),
+ },
+ {
+ label: translate("browse:dagBundles.columns.version"),
+ value: (
+ <DagBundleVersion
+ bundleUrl={bundle.bundle_url}
+ lastRefreshed={bundle.last_refreshed}
+ version={bundle.version}
+ />
+ ),
+ },
+ {
+ label: translate("browse:dagBundles.columns.lastRefreshed"),
+ value:
+ bundle.last_refreshed === null ? (
+ translate("browse:dagBundles.neverRefreshed")
+ ) : (
+ <Time datetime={bundle.last_refreshed} />
+ ),
+ },
+ { label: translate("common:dag_other"), value: bundle.dag_count },
+ {
+ label: translate("browse:dagBundles.columns.importErrors"),
+ value: <ImportErrorCount count={bundle.import_error_count} />,
+ },
+ ...(showTeam
+ ? [
+ {
+ label: translate("common:dagDetails.team"),
+ value: <TeamName teamName={bundle.team_name} />,
+ },
+ ]
+ : []),
+ ];
+
+ return <HeaderCard icon={<FiPackage />} stats={stats} title={bundle.name}
type="dagBundle" />;
+};
diff --git a/airflow-core/src/airflow/ui/src/pages/DagBundle/index.tsx
b/airflow-core/src/airflow/ui/src/pages/DagBundle/index.tsx
new file mode 100644
index 00000000000..87137650a4d
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/pages/DagBundle/index.tsx
@@ -0,0 +1,19 @@
+/*!
+ * 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.
+ */
+export { DagBundle } from "./DagBundle";
diff --git a/airflow-core/src/airflow/ui/src/pages/DagBundles/DagBundles.tsx
b/airflow-core/src/airflow/ui/src/pages/DagBundles/DagBundles.tsx
index 39603fbce60..8a4475c24a3 100644
--- a/airflow-core/src/airflow/ui/src/pages/DagBundles/DagBundles.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/DagBundles/DagBundles.tsx
@@ -16,7 +16,7 @@
* specific language governing permissions and limitations
* under the License.
*/
-import { Badge, Box, Link, Text } from "@chakra-ui/react";
+import { Badge, Box, Text } from "@chakra-ui/react";
import type { ColumnDef } from "@tanstack/react-table";
import type { TFunction } from "i18next";
import { useTranslation } from "react-i18next";
@@ -24,24 +24,19 @@ import { useTranslation } from "react-i18next";
import { useDagBundleServiceGetDagBundles } from "openapi/queries";
import type { DagBundleResponse } from "openapi/requests/types.gen";
-import { Tooltip } from "src/system-components";
+import { RouterLink, Tooltip } from "src/system-components";
+import { DagBundleVersion } from "src/components/DagBundleVersion";
import { DataTable } from "src/components/DataTable";
import { useTableURLState } from "src/components/DataTable/useTableUrlState";
import { ErrorAlert } from "src/components/ErrorAlert";
+import { ImportErrorCount } from "src/components/ImportErrorCount";
import { TeamName } from "src/components/TeamName";
import Time from "src/components/Time";
import { useConfig } from "src/queries/useConfig";
+import { useDagBundleRefetchInterval } from
"src/queries/useDagBundleRefetchInterval";
import { type DurationFormat, useDocumentTitle, useDurationFormat } from
"src/utils";
-import { useAutoRefresh } from "src/utils/query";
-
-// A bundle row changes at most once per the bundle's `refresh_interval`
(default 300s), so
-// `[api] auto_refresh_interval` (default 3s, tuned for Grid/Graph run state)
is far too fast here.
-// Floor it rather than ignore it, so turning auto-refresh off still turns
this page off.
-const MIN_REFETCH_INTERVAL_MS = 10_000;
-
-const SHORT_VERSION_LENGTH = 7;
type BundleRow = { row: { original: DagBundleResponse } };
@@ -52,6 +47,11 @@ const createColumns = (
): Array<ColumnDef<DagBundleResponse>> => [
{
accessorKey: "name",
+ cell: ({ row: { original } }: BundleRow) => (
+ <RouterLink fontWeight="bold"
to={`/dag_bundles/${encodeURIComponent(original.name)}`}>
+ {original.name}
+ </RouterLink>
+ ),
header: translate("browse:dagBundles.columns.name"),
},
...(multiTeam
@@ -66,41 +66,13 @@ const createColumns = (
: []),
{
accessorKey: "version",
- cell: ({ row: { original } }: BundleRow) => {
- if (original.version === null) {
- // A versioning-capable bundle also reports null until its first
successful refresh, so an
- // absent last_refreshed is what separates "not refreshed yet" from
"not versioned".
- return (
- <Text color="fg.muted">
- {original.last_refreshed === null
- ? translate("browse:dagBundles.notRefreshedYet")
- : translate("browse:dagBundles.notVersioned")}
- </Text>
- );
- }
-
- // A git bundle stores the full 40-char hexsha, so show the prefix an
author recognises and
- // keep the whole value on hover.
- const short = original.version.slice(0, SHORT_VERSION_LENGTH);
-
- return (
- <Tooltip content={original.version}>
- {original.bundle_url === null ? (
- <Text fontFamily="mono">{short}</Text>
- ) : (
- <Link
- color="fg.info"
- fontFamily="mono"
- href={original.bundle_url}
- rel="noreferrer"
- target="_blank"
- >
- {short}
- </Link>
- )}
- </Tooltip>
- );
- },
+ cell: ({ row: { original } }: BundleRow) => (
+ <DagBundleVersion
+ bundleUrl={original.bundle_url}
+ lastRefreshed={original.last_refreshed}
+ version={original.version}
+ />
+ ),
header: translate("browse:dagBundles.columns.version"),
},
{
@@ -117,20 +89,7 @@ const createColumns = (
},
{
accessorKey: "import_error_count",
- cell: ({ row: { original } }: BundleRow) => {
- // null means the user may not read import errors, which is not the same
as "none".
- if (original.import_error_count === null) {
- return <Text color="fg.muted">-</Text>;
- }
-
- return original.import_error_count === 0 ? (
- <Text color="fg.muted">0</Text>
- ) : (
- <Badge colorPalette="failed" variant="solid">
- {original.import_error_count}
- </Badge>
- );
- },
+ cell: ({ row: { original } }: BundleRow) => <ImportErrorCount
count={original.import_error_count} />,
enableSorting: false,
header: translate("browse:dagBundles.columns.importErrors"),
},
@@ -160,13 +119,7 @@ export const DagBundles = () => {
const [sort] = sorting;
const orderBy = sort ? [`${sort.desc ? "-" : ""}${sort.id}`] : undefined;
- // `useAutoRefresh` with no dagId reduces to `[api] auto_refresh_interval`,
where 0 is how an
- // operator turns auto-refresh off -- so it has to short-circuit before the
floor below.
- const configuredInterval = useAutoRefresh({});
- const refetchInterval =
- configuredInterval === false || configuredInterval === 0
- ? false
- : Math.max(configuredInterval, MIN_REFETCH_INTERVAL_MS);
+ const refetchInterval = useDagBundleRefetchInterval();
const { data, error, isFetching, isLoading } =
useDagBundleServiceGetDagBundles(
{
@@ -191,7 +144,7 @@ export const DagBundles = () => {
initialState={tableURLState}
isFetching={isFetching}
isLoading={isLoading}
- modelName="browse:dagBundles.bundle"
+ modelName="common:dagBundle"
onStateChange={setTableURLState}
total={data?.total_entries}
/>
diff --git
a/airflow-core/src/airflow/ui/src/queries/useDagBundleRefetchInterval.ts
b/airflow-core/src/airflow/ui/src/queries/useDagBundleRefetchInterval.ts
new file mode 100644
index 00000000000..62e9cd38bf3
--- /dev/null
+++ b/airflow-core/src/airflow/ui/src/queries/useDagBundleRefetchInterval.ts
@@ -0,0 +1,38 @@
+/*!
+ * 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.
+ */
+import { useAutoRefresh } from "src/utils/query";
+
+// A bundle row changes at most once per the bundle's `refresh_interval`
(default 300s), so
+// `[api] auto_refresh_interval` (default 3s, tuned for Grid/Graph run state)
is far too fast here.
+// Floor it rather than ignore it, so turning auto-refresh off still turns
these pages off.
+const MIN_REFETCH_INTERVAL_MS = 10_000;
+
+/**
+ * How often the Dag bundle pages should poll, or false when polling is off.
+ *
+ * `useAutoRefresh` with no dagId reduces to `[api] auto_refresh_interval`,
where 0 is how an
+ * operator turns auto-refresh off -- so it has to short-circuit before the
floor.
+ */
+export const useDagBundleRefetchInterval = (): number | false => {
+ const configuredInterval = useAutoRefresh({});
+
+ return configuredInterval === false || configuredInterval === 0
+ ? false
+ : Math.max(configuredInterval, MIN_REFETCH_INTERVAL_MS);
+};
diff --git a/airflow-core/src/airflow/ui/src/router.tsx
b/airflow-core/src/airflow/ui/src/router.tsx
index 8f5ee4e232e..9b62505c6a9 100644
--- a/airflow-core/src/airflow/ui/src/router.tsx
+++ b/airflow-core/src/airflow/ui/src/router.tsx
@@ -38,6 +38,7 @@ import { Code } from "src/pages/Dag/Code";
import { Details as DagDetails } from "src/pages/Dag/Details";
import { Overview } from "src/pages/Dag/Overview";
import { Tasks } from "src/pages/Dag/Tasks";
+import { DagBundle } from "src/pages/DagBundle";
import { DagBundles } from "src/pages/DagBundles";
import { DagRuns } from "src/pages/DagRuns";
import { DagsList } from "src/pages/DagsList";
@@ -163,6 +164,10 @@ export const routerConfig = [
element: <DagBundles />,
path: "dag_bundles",
},
+ {
+ element: <DagBundle />,
+ path: "dag_bundles/:bundleName",
+ },
{
element: <Deadlines />,
path: "deadlines",
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_bundles.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_bundles.py
index 7c1ff834675..faee8234661 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_bundles.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_bundles.py
@@ -62,11 +62,14 @@ UNREADABLE_FILE = "secret.py"
GIT_VERSION = "8f0e5b1c9a2d4e6f8a0b1c2d3e4f5a6b7c8d9e0f"
REFRESHED_AT = datetime(2026, 9, 10, 12, 0, tzinfo=timezone.utc)
+PARSED_AT = datetime(2026, 9, 10, 12, 1, tzinfo=timezone.utc)
+PARSE_DURATION = 0.125
GONE_REFRESHED_AT = datetime(2026, 9, 1, 8, 30, tzinfo=timezone.utc)
WITH_DAGS = (GIT_BUNDLE, LOCAL_BUNDLE, GONE_BUNDLE, OTHER_TEAM_BUNDLE)
# Everything except the Dag in OTHER_TEAM_BUNDLE.
-READABLE_DAG_IDS = {f"dag_in_{name}" for name in (GIT_BUNDLE, LOCAL_BUNDLE,
GONE_BUNDLE)}
+REMOVED_DAG_ID = "removed_from_dag_0"
+READABLE_DAG_IDS = {f"dag_in_{name}" for name in (GIT_BUNDLE, LOCAL_BUNDLE,
GONE_BUNDLE)} | {REMOVED_DAG_ID}
# Sorted by name, which is the endpoint's default order.
READABLE_BUNDLES = [GIT_BUNDLE, GONE_BUNDLE, LOCAL_BUNDLE]
@@ -149,8 +152,27 @@ def bundles() -> Generator[None, None, None]:
relative_fileloc=f"dag_{index}.py",
bundle_name=bundle_name,
is_paused=False,
+ # ``is_stale`` defaults to True on the model, and a parsed
Dag is not stale.
+ is_stale=False,
+ last_parsed_time=PARSED_AT,
+ last_parse_duration=PARSE_DURATION,
)
)
+ # Removed from REGISTERED_FILE but never deleted: stale, and frozen at
an older parse
+ # that took far longer. An unrestricted MAX over the file would pair
this duration with
+ # the live row's newer timestamp.
+ session.add(
+ DagModel(
+ dag_id=REMOVED_DAG_ID,
+ fileloc=REGISTERED_FILE,
+ relative_fileloc=REGISTERED_FILE,
+ bundle_name=GIT_BUNDLE,
+ is_paused=False,
+ is_stale=True,
+ last_parsed_time=PARSED_AT - timedelta(hours=1),
+ last_parse_duration=PARSE_DURATION * 100,
+ )
+ )
# A co-located Dag in a visible bundle that the caller cannot read --
deliberately absent
# from READABLE_DAG_IDS.
session.add(
@@ -160,6 +182,7 @@ def bundles() -> Generator[None, None, None]:
relative_fileloc=UNREADABLE_FILE,
bundle_name=GIT_BUNDLE,
is_paused=False,
+ is_stale=False,
)
)
session.add_all(
@@ -682,3 +705,244 @@ class TestDaglessBundleTeamScoping:
names = [bundle["name"] for bundle in body["dag_bundles"]]
assert names == READABLE_BUNDLES
assert body["total_entries"] == len(READABLE_BUNDLES)
+
+
+class TestGetDagBundle:
+ def test_should_raise_401_unauthenticated(self,
unauthenticated_test_client):
+ assert
unauthenticated_test_client.get(f"/dagBundles/{GIT_BUNDLE}").status_code == 401
+
+ def test_should_raise_403_unauthorized(self, unauthorized_test_client):
+ assert
unauthorized_test_client.get(f"/dagBundles/{GIT_BUNDLE}").status_code == 403
+
+ def test_returns_the_bundle(self, dag_scoped_client):
+ response = dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}")
+
+ assert response.status_code == 200
+ body = response.json()
+ assert body["name"] == GIT_BUNDLE
+ assert body["version"] == GIT_VERSION
+ assert body["last_refreshed"] == "2026-09-10T12:00:00Z"
+ assert body["active"] is True
+ assert body["bundle_url"] ==
f"https://github.com/example/repo/tree/{GIT_VERSION}/dags"
+
+ def test_dag_count_excludes_dags_the_caller_may_not_read(self,
dag_scoped_client):
+ """``GIT_BUNDLE`` holds two Dags, one of them absent from
``READABLE_DAG_IDS``."""
+ body = dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}").json()
+
+ assert body["dag_count"] == 1
+
+ def test_404_for_a_bundle_whose_dags_are_not_readable(self,
dag_scoped_client):
+ """A 404 rather than a 403, so the response does not confirm the
bundle exists."""
+ assert
dag_scoped_client.get(f"/dagBundles/{OTHER_TEAM_BUNDLE}").status_code == 404
+
+ def test_404_for_a_bundle_with_no_dags_without_the_admin_view(self,
dag_scoped_client):
+ """``dag_scoped_client`` is denied ``IMPORT_ERRORS_ALL``, so Dag
scoping alone decides."""
+ assert
dag_scoped_client.get(f"/dagBundles/{DAGLESS_BUNDLE}").status_code == 404
+ assert
dag_scoped_client.get(f"/dagBundles/{DAGLESS_BUNDLE}/files").status_code == 404
+
+ def
test_shows_a_bundle_with_no_dags_to_a_caller_holding_the_admin_view(self,
admin_client):
+ """
+ The first-deploy-is-broken case: a bundle whose only file fails before
defining a Dag.
+
+ It has no Dag to authorize against, so it rides on
``IMPORT_ERRORS_ALL`` the way the
+ collection route does; the detail route must not re-derive a stricter
rule and hide it.
+ """
+ response = admin_client.get(f"/dagBundles/{DAGLESS_BUNDLE}")
+
+ assert response.status_code == 200
+ assert response.json()["dag_count"] == 0
+ assert
admin_client.get(f"/dagBundles/{DAGLESS_BUNDLE}/files").status_code == 200
+
+ def test_404_for_an_unknown_bundle(self, dag_scoped_client):
+ assert dag_scoped_client.get("/dagBundles/no_such_bundle").status_code
== 404
+
+ def test_import_error_count_is_gated_like_the_collection(self,
admin_client, viewer_client):
+ """
+ The admin sees the unregistered-file error as well; the viewer sees
only the registered one.
+
+ Same two-part authorization as ``GET /importErrors``, asserted here so
the detail route
+ cannot drift from the collection route it shares a helper with.
+ """
+ assert
admin_client.get(f"/dagBundles/{GIT_BUNDLE}").json()["import_error_count"] == 2
+ assert
viewer_client.get(f"/dagBundles/{GIT_BUNDLE}").json()["import_error_count"] == 1
+
+ def test_import_error_count_is_withheld_without_permission(self,
admin_client):
+ auth_manager = admin_client.app.state.auth_manager
+ with mock.patch.object(auth_manager, "authorize_view", autospec=True,
return_value=False):
+ body = admin_client.get(f"/dagBundles/{GIT_BUNDLE}").json()
+
+ assert body["import_error_count"] is None
+
+ @conf_vars({("core", "multi_team"): "True"})
+ def test_reports_the_owning_team_in_multi_team_mode(self,
dag_scoped_client, session):
+ """Without ``multi_team`` on -- the shipped default -- the team lookup
is skipped entirely."""
+ session.add(Team(name=TEAM_NAME))
+ session.commit()
+ session.execute(
+
insert(dag_bundle_team_association_table).values(dag_bundle_name=GIT_BUNDLE,
team_name=TEAM_NAME)
+ )
+ session.commit()
+
+ assert
dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}").json()["team_name"] ==
TEAM_NAME
+
+ def test_bundle_url_is_withheld_without_dag_version_read(self,
dag_scoped_client):
+ """
+ A rendered bundle url is otherwise only reachable through ``GET
/dags/{dag_id}/dagVersions``.
+
+ That route additionally requires Dag *version* read, so a role that
grants Dag read and
+ withholds version read must not get the repository address here
instead.
+ """
+ auth_manager = dag_scoped_client.app.state.auth_manager
+ real = auth_manager.is_authorized_dag
+
+ def _deny_versions(*args, **kwargs):
+ if kwargs.get("access_entity") is DagAccessEntity.VERSION:
+ return False
+ return real(*args, **kwargs)
+
+ with mock.patch.object(auth_manager, "is_authorized_dag",
autospec=True, side_effect=_deny_versions):
+ body = dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}").json()
+
+ assert body["version"] == GIT_VERSION
+ assert body["bundle_url"] is None
+
+ def test_dag_count_excludes_stale_dags(self, dag_scoped_client, session):
+ """``GIT_BUNDLE``'s Dags all go stale, so nothing live is left to
count."""
+ session.execute(update(DagModel).where(DagModel.bundle_name ==
GIT_BUNDLE).values(is_stale=True))
+ session.commit()
+
+ body = dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}").json()
+
+ assert body["dag_count"] == 0
+
+
+class TestGetDagBundleFiles:
+ def test_should_raise_401_unauthenticated(self,
unauthenticated_test_client):
+ assert
unauthenticated_test_client.get(f"/dagBundles/{GIT_BUNDLE}/files").status_code
== 401
+
+ def test_should_raise_403_unauthorized(self, unauthorized_test_client):
+ assert
unauthorized_test_client.get(f"/dagBundles/{GIT_BUNDLE}/files").status_code ==
403
+
+ def test_lists_files_ordered_by_path(self, admin_client):
+ response = admin_client.get(f"/dagBundles/{GIT_BUNDLE}/files")
+
+ assert response.status_code == 200
+ body = response.json()
+ assert [file["relative_fileloc"] for file in body["dag_bundle_files"]]
== [
+ UNREGISTERED_FILE,
+ REGISTERED_FILE,
+ ]
+ assert body["total_entries"] == 2
+
+ def test_reports_parse_time_and_duration(self, dag_scoped_client):
+ body = dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}/files").json()
+
+ by_path = {file["relative_fileloc"]: file for file in
body["dag_bundle_files"]}
+ assert by_path[REGISTERED_FILE]["dag_count"] == 1
+ assert by_path[REGISTERED_FILE]["last_parsed_time"] ==
"2026-09-10T12:01:00Z"
+ assert by_path[REGISTERED_FILE]["last_parse_duration"] ==
PARSE_DURATION
+
+ def test_a_file_that_registered_no_dag_is_listed_with_its_error(self,
admin_client):
+ """
+ The case the page most needs to show: a file that failed before
defining a Dag.
+
+ It has no Dag to authorize on, so it is admin-gated, and without it a
brand new broken
+ file would be invisible on the very page someone opens to find out why.
+ """
+ body = admin_client.get(f"/dagBundles/{GIT_BUNDLE}/files").json()
+
+ by_path = {file["relative_fileloc"]: file for file in
body["dag_bundle_files"]}
+ assert by_path[UNREGISTERED_FILE]["dag_count"] == 0
+ assert by_path[UNREGISTERED_FILE]["import_error_count"] == 1
+ assert by_path[UNREGISTERED_FILE]["last_parsed_time"] is None
+ assert by_path[UNREGISTERED_FILE]["last_parse_duration"] is None
+
+ def test_excludes_a_file_whose_dag_is_not_readable(self, admin_client,
viewer_client):
+ """``UNREADABLE_FILE`` is registered, so only the readable-Dag filter
keeps it out."""
+ for client in (admin_client, viewer_client):
+ body = client.get(f"/dagBundles/{GIT_BUNDLE}/files").json()
+ assert UNREADABLE_FILE not in {file["relative_fileloc"] for file
in body["dag_bundle_files"]}
+
+ def test_viewer_does_not_see_the_unregistered_file(self, viewer_client):
+ body = viewer_client.get(f"/dagBundles/{GIT_BUNDLE}/files").json()
+
+ assert [file["relative_fileloc"] for file in body["dag_bundle_files"]]
== [REGISTERED_FILE]
+ assert body["total_entries"] == 1
+
+ def test_import_error_count_is_withheld_without_permission(self,
admin_client):
+ """``None``, not 0: "you may not see this" must not read as "nothing
is wrong"."""
+ auth_manager = admin_client.app.state.auth_manager
+ with mock.patch.object(auth_manager, "authorize_view", autospec=True,
return_value=False):
+ body = admin_client.get(f"/dagBundles/{GIT_BUNDLE}/files").json()
+
+ assert [file["relative_fileloc"] for file in body["dag_bundle_files"]]
== [REGISTERED_FILE]
+ assert body["dag_bundle_files"][0]["import_error_count"] is None
+
+ def test_a_file_whose_dags_went_stale_stays_listed_while_it_has_an_error(
+ self, dag_scoped_client, session
+ ):
+ """Marking a file's Dags stale must not hide the row that explains the
breakage."""
+ session.execute(
+ update(DagModel)
+ .where(DagModel.relative_fileloc == REGISTERED_FILE,
DagModel.bundle_name == GIT_BUNDLE)
+ .values(is_stale=True)
+ )
+ session.commit()
+
+ body = dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}/files").json()
+
+ by_path = {file["relative_fileloc"]: file for file in
body["dag_bundle_files"]}
+ assert by_path[REGISTERED_FILE]["dag_count"] == 0
+ assert by_path[REGISTERED_FILE]["import_error_count"] == 1
+
+ def test_parse_duration_ignores_a_stale_dag_from_an_older_parse(self,
dag_scoped_client):
+ """
+ A Dag removed from a file keeps its row, frozen at the parse it was
last seen in.
+
+ Aggregating the two columns independently would pair that older,
slower duration with the
+ newer timestamp and report a parse time that never happened.
+ """
+ body = dag_scoped_client.get(f"/dagBundles/{GIT_BUNDLE}/files").json()
+
+ by_path = {file["relative_fileloc"]: file for file in
body["dag_bundle_files"]}
+ assert by_path[REGISTERED_FILE]["last_parsed_time"] ==
"2026-09-10T12:01:00Z"
+ assert by_path[REGISTERED_FILE]["last_parse_duration"] ==
PARSE_DURATION
+
+ def test_a_file_with_no_live_dag_and_no_error_is_dropped(self,
dag_scoped_client, session):
+ """A file whose Dags are all stale and which has no error is one the
bundle no longer has."""
+ session.execute(update(DagModel).where(DagModel.bundle_name ==
LOCAL_BUNDLE).values(is_stale=True))
+ session.commit()
+
+ body =
dag_scoped_client.get(f"/dagBundles/{LOCAL_BUNDLE}/files").json()
+
+ assert body["dag_bundle_files"] == []
+ assert body["total_entries"] == 0
+
+ def test_a_file_with_no_error_reports_zero(self, dag_scoped_client):
+ body =
dag_scoped_client.get(f"/dagBundles/{LOCAL_BUNDLE}/files").json()
+
+ assert [file["import_error_count"] for file in
body["dag_bundle_files"]] == [0]
+
+ @pytest.mark.parametrize(
+ ("params", "expected"),
+ [
+ pytest.param({"limit": 1}, [UNREGISTERED_FILE], id="limit"),
+ pytest.param({"offset": 1}, [REGISTERED_FILE], id="offset"),
+ pytest.param({"limit": 1, "offset": 1}, [REGISTERED_FILE],
id="limit-and-offset"),
+ # ``limit`` is a non-negative int, and ``Select.limit(0)`` returns
nothing, so zero
+ # has to mean nothing here too rather than being read as
"unlimited".
+ pytest.param({"limit": 0}, [], id="limit-zero"),
+ ],
+ )
+ def test_pagination(self, admin_client, params, expected):
+ body = admin_client.get(f"/dagBundles/{GIT_BUNDLE}/files",
params=params).json()
+
+ assert [file["relative_fileloc"] for file in body["dag_bundle_files"]]
== expected
+ # The total is the whole file list, not the page.
+ assert body["total_entries"] == 2
+
+ def test_404_for_a_bundle_whose_dags_are_not_readable(self,
dag_scoped_client):
+ assert
dag_scoped_client.get(f"/dagBundles/{OTHER_TEAM_BUNDLE}/files").status_code ==
404
+
+ def test_404_for_an_unknown_bundle(self, dag_scoped_client):
+ assert
dag_scoped_client.get("/dagBundles/no_such_bundle/files").status_code == 404
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index cd3210e6699..23329487c61 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -493,6 +493,95 @@ class DAGTagCollectionResponse(BaseModel):
total_entries: Annotated[int, Field(title="Total Entries")]
+class DagBundleDetailResponse(BaseModel):
+ """
+ Dag bundle serializer for the single-bundle response.
+ """
+
+ name: Annotated[str, Field(title="Name")]
+ active: Annotated[
+ bool | None,
+ Field(
+ description="Whether the bundle is still present in this
deployment's configuration.",
+ title="Active",
+ ),
+ ]
+ version: Annotated[
+ str | None,
+ Field(
+ description="The latest version Airflow has seen for the bundle.
Null when the bundle does not support versioning, or when no Dag processor has
refreshed it successfully yet.",
+ title="Version",
+ ),
+ ]
+ last_refreshed: Annotated[
+ datetime | None,
+ Field(
+ description="When a Dag processor last successfully refreshed the
bundle. It advances even when the version did not change, and a failed refresh
leaves it untouched.",
+ title="Last Refreshed",
+ ),
+ ]
+ bundle_url: Annotated[
+ str | None,
+ Field(
+ description="A link to view the bundle at ``version``, when one is
configured and the caller may read Dag versions.",
+ title="Bundle Url",
+ ),
+ ]
+ team_name: Annotated[
+ str | None,
+ Field(description="The team owning the bundle, in a multi-team
deployment.", title="Team Name"),
+ ]
+ import_error_count: Annotated[
+ int | None,
+ Field(
+ description="Number of Dag import errors recorded against this
bundle that the caller is permitted to see, counted on the same terms as ``GET
/importErrors``. Null when the caller may not read import errors.",
+ title="Import Error Count",
+ ),
+ ]
+ dag_count: Annotated[
+ int,
+ Field(
+ description="Number of live Dags recorded against the bundle that
the caller is permitted to see, counted on the same terms as ``GET /dags``.",
+ title="Dag Count",
+ ),
+ ]
+
+
+class DagBundleFileResponse(BaseModel):
+ """
+ A file in a Dag bundle, as the Dag processor last saw it.
+ """
+
+ relative_fileloc: Annotated[str, Field(title="Relative Fileloc")]
+ dag_count: Annotated[
+ int,
+ Field(
+ description="Number of live Dags the file defines that the caller
may read.", title="Dag Count"
+ ),
+ ]
+ last_parsed_time: Annotated[
+ datetime | None,
+ Field(
+ description="When the file was last parsed, or null if it has
never parsed successfully.",
+ title="Last Parsed Time",
+ ),
+ ]
+ last_parse_duration: Annotated[
+ float | None,
+ Field(
+ description="How long the last successful parse of the file took,
in seconds.",
+ title="Last Parse Duration",
+ ),
+ ]
+ import_error_count: Annotated[
+ int | None,
+ Field(
+ description="Number of import errors recorded against the file,
which is at most one. Null when the caller may not read import errors --
deliberately not zero, which would read as a file with nothing wrong.",
+ title="Import Error Count",
+ ),
+ ]
+
+
class DagBundleResponse(BaseModel):
"""
Dag bundle serializer for responses.
@@ -2022,6 +2111,15 @@ class DagBundleCollectionResponse(BaseModel):
total_entries: Annotated[int, Field(title="Total Entries")]
+class DagBundleFileCollectionResponse(BaseModel):
+ """
+ Dag bundle file collection response.
+ """
+
+ dag_bundle_files: Annotated[list[DagBundleFileResponse], Field(title="Dag
Bundle Files")]
+ total_entries: Annotated[int, Field(title="Total Entries")]
+
+
class DagProcessorInfoResponse(BaseModel):
"""
DagProcessor info serializer for responses.
diff --git a/providers/fab/docs/auth-manager/access-control.rst
b/providers/fab/docs/auth-manager/access-control.rst
index c7cad5fd947..ffb1849aae9 100644
--- a/providers/fab/docs/auth-manager/access-control.rst
+++ b/providers/fab/docs/auth-manager/access-control.rst
@@ -360,6 +360,14 @@ Stable API Permissions
- GET
- DAGs.can_read
- Viewer
+ * - ``/api/v2/dagBundles/{bundle_name}``
+ - GET
+ - DAGs.can_read
+ - Viewer
+ * - ``/api/v2/dagBundles/{bundle_name}/files``
+ - GET
+ - DAGs.can_read
+ - Viewer
* - ``/api/v2/dagSources/{dag_id}``
- GET
- DAGs.can_read, DAG Code.can_read