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

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 

18 

19from typing import Annotated 

20 

21from fastapi import Depends, HTTPException, status 

22from sqlalchemy import select 

23from sqlalchemy.orm import joinedload 

24 

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 

46 

47dag_versions_router = AirflowRouter(tags=["DagVersion"], prefix="/dags/{dag_id}/dagVersions") 

48 

49 

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 ) 

70 

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 ) 

76 

77 return dag_version 

78 

79 

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. 

112 

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)) 

116 

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) 

120 

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) 

130 

131 return DAGVersionCollectionResponse( 

132 dag_versions=dag_versions, 

133 total_entries=total_entries, 

134 )