Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/dag_versions.py: 97%
30 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
1# Licensed to the Apache Software Foundation (ASF) under one
2# or more contributor license agreements. See the NOTICE file
3# distributed with this work for additional information
4# regarding copyright ownership. The ASF licenses this file
5# to you under the Apache License, Version 2.0 (the
6# "License"); you may not use this file except in compliance
7# with the License. You may obtain a copy of the License at
8#
9# http://www.apache.org/licenses/LICENSE-2.0
10#
11# Unless required by applicable law or agreed to in writing,
12# software distributed under the License is distributed on an
13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14# KIND, either express or implied. See the License for the
15# specific language governing permissions and limitations
16# under the License.
17from __future__ import annotations
19from typing import Annotated
21from fastapi import Depends, HTTPException, status
22from sqlalchemy import select
23from sqlalchemy.orm import joinedload
25from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity
26from airflow.api_fastapi.common.dagbag import DagBagDep, get_latest_version_of_dag
27from airflow.api_fastapi.common.db.common import SessionDep, paginated_select
28from airflow.api_fastapi.common.parameters import (
29 FilterParam,
30 QueryLimit,
31 QueryOffset,
32 SortParam,
33 filter_param_factory,
34)
35from airflow.api_fastapi.common.router import AirflowRouter
36from airflow.api_fastapi.core_api.datamodels.dag_versions import (
37 DAGVersionCollectionResponse,
38 DagVersionResponse,
39)
40from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
41from airflow.api_fastapi.core_api.security import (
42 ReadableDagVersionsFilterDep,
43 requires_access_dag,
44)
45from airflow.models.dag_version import DagVersion
47dag_versions_router = AirflowRouter(tags=["DagVersion"], prefix="/dags/{dag_id}/dagVersions")
50@dag_versions_router.get(
51 "/{version_number}",
52 responses=create_openapi_http_exception_doc(
53 [
54 status.HTTP_404_NOT_FOUND,
55 ]
56 ),
57 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.VERSION))],
58)
59def get_dag_version(
60 dag_id: str,
61 version_number: int,
62 session: SessionDep,
63) -> DagVersionResponse:
64 """Get one Dag Version."""
65 dag_version = session.scalar(
66 select(DagVersion)
67 .filter_by(dag_id=dag_id, version_number=version_number)
68 .options(joinedload(DagVersion.dag_model))
69 )
71 if dag_version is None:
72 raise HTTPException(
73 status.HTTP_404_NOT_FOUND,
74 f"The DagVersion with dag_id: `{dag_id}` and version_number: `{version_number}` was not found",
75 )
77 return dag_version
80@dag_versions_router.get(
81 "",
82 responses=create_openapi_http_exception_doc(
83 [
84 status.HTTP_404_NOT_FOUND,
85 ],
86 ),
87 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.VERSION))],
88)
89def get_dag_versions(
90 dag_id: str,
91 session: SessionDep,
92 limit: QueryLimit,
93 offset: QueryOffset,
94 version_number: Annotated[
95 FilterParam[int], Depends(filter_param_factory(DagVersion.version_number, int))
96 ],
97 bundle_name: Annotated[FilterParam[str], Depends(filter_param_factory(DagVersion.bundle_name, str))],
98 bundle_version: Annotated[
99 FilterParam[str | None], Depends(filter_param_factory(DagVersion.bundle_version, str | None))
100 ],
101 order_by: Annotated[
102 SortParam,
103 Depends(
104 SortParam(["id", "version_number", "bundle_name", "bundle_version"], DagVersion).dynamic_depends()
105 ),
106 ],
107 dag_bag: DagBagDep,
108 readable_dag_versions_filter: ReadableDagVersionsFilterDep,
109) -> DAGVersionCollectionResponse:
110 """
111 Get all Dag Versions.
113 This endpoint allows specifying `~` as the dag_id to retrieve Dag Versions for all Dags.
114 """
115 query = select(DagVersion).options(joinedload(DagVersion.dag_model), joinedload(DagVersion.bundle))
117 if dag_id != "~": 117 ↛ 121line 117 didn't jump to line 121 because the condition on line 117 was always true
118 get_latest_version_of_dag(dag_bag, dag_id, session)
119 query = query.filter(DagVersion.dag_id == dag_id)
121 dag_versions_select, total_entries = paginated_select(
122 statement=query,
123 filters=[version_number, bundle_name, bundle_version, readable_dag_versions_filter],
124 order_by=order_by,
125 offset=offset,
126 limit=limit,
127 session=session,
128 )
129 dag_versions = session.scalars(dag_versions_select)
131 return DAGVersionCollectionResponse(
132 dag_versions=dag_versions,
133 total_entries=total_entries,
134 )