Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/variables.py: 97%

51 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, Query, status 

22from sqlalchemy import delete, select 

23 

24from airflow.api_fastapi.common.db.common import SessionDep, paginated_select 

25from airflow.api_fastapi.common.parameters import ( 

26 QueryLimit, 

27 QueryOffset, 

28 QueryVariableKeyPatternSearch, 

29 QueryVariableKeyPrefixPatternSearch, 

30 SortParam, 

31) 

32from airflow.api_fastapi.common.router import AirflowRouter 

33from airflow.api_fastapi.core_api.datamodels.common import BulkBody, BulkResponse 

34from airflow.api_fastapi.core_api.datamodels.variables import ( 

35 VariableBody, 

36 VariableCollectionResponse, 

37 VariableResponse, 

38) 

39from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc 

40from airflow.api_fastapi.core_api.security import ( 

41 ReadableVariablesFilterDep, 

42 requires_access_variable, 

43 requires_access_variable_bulk, 

44) 

45from airflow.api_fastapi.core_api.services.public.variables import ( 

46 BulkVariableService, 

47 update_orm_from_pydantic, 

48) 

49from airflow.api_fastapi.logging.decorators import action_logging 

50from airflow.models.variable import Variable 

51 

52variables_router = AirflowRouter(tags=["Variable"], prefix="/variables") 

53 

54 

55@variables_router.delete( 

56 "/{variable_key:path}", 

57 status_code=status.HTTP_204_NO_CONTENT, 

58 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

59 dependencies=[Depends(action_logging()), Depends(requires_access_variable("DELETE"))], 

60) 

61def delete_variable( 

62 variable_key: str, 

63 session: SessionDep, 

64): 

65 """Delete a variable entry.""" 

66 # Like the other endpoints (get, patch), we do not use Variable.delete/get/set here because these methods 

67 # are intended to be used in task execution environment (execution API) 

68 result = session.execute(delete(Variable).where(Variable.key == variable_key)) 

69 rows = getattr(result, "rowcount", 0) or 0 

70 if rows == 0: 

71 raise HTTPException( 

72 status.HTTP_404_NOT_FOUND, f"The Variable with key: `{variable_key}` was not found" 

73 ) 

74 

75 

76@variables_router.get( 

77 "/{variable_key:path}", 

78 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

79 dependencies=[Depends(requires_access_variable("GET"))], 

80) 

81def get_variable( 

82 variable_key: str, 

83 session: SessionDep, 

84) -> VariableResponse: 

85 """Get a variable entry.""" 

86 variable = session.scalar(select(Variable).where(Variable.key == variable_key).limit(1)) 

87 

88 if variable is None: 

89 raise HTTPException( 

90 status.HTTP_404_NOT_FOUND, f"The Variable with key: `{variable_key}` was not found" 

91 ) 

92 

93 return variable 

94 

95 

96@variables_router.get( 

97 "", 

98 dependencies=[Depends(requires_access_variable("GET"))], 

99) 

100def get_variables( 

101 limit: QueryLimit, 

102 offset: QueryOffset, 

103 order_by: Annotated[ 

104 SortParam, 

105 Depends( 

106 SortParam( 

107 ["key", "id", "_val", "description", "is_encrypted", "team_name"], 

108 Variable, 

109 ).dynamic_depends() 

110 ), 

111 ], 

112 readable_variables_filter: ReadableVariablesFilterDep, 

113 session: SessionDep, 

114 variable_key_pattern: QueryVariableKeyPatternSearch, 

115 variable_key_prefix_pattern: QueryVariableKeyPrefixPatternSearch, 

116) -> VariableCollectionResponse: 

117 """Get all Variables entries.""" 

118 variable_select, total_entries = paginated_select( 

119 statement=select(Variable), 

120 filters=[variable_key_pattern, variable_key_prefix_pattern, readable_variables_filter], 

121 order_by=order_by, 

122 offset=offset, 

123 limit=limit, 

124 session=session, 

125 ) 

126 

127 variables = session.scalars(variable_select) 

128 

129 return VariableCollectionResponse( 

130 variables=variables, 

131 total_entries=total_entries, 

132 ) 

133 

134 

135@variables_router.patch( 

136 "/{variable_key:path}", 

137 responses=create_openapi_http_exception_doc( 

138 [ 

139 status.HTTP_400_BAD_REQUEST, 

140 status.HTTP_404_NOT_FOUND, 

141 ] 

142 ), 

143 dependencies=[Depends(action_logging()), Depends(requires_access_variable("PUT"))], 

144) 

145def patch_variable( 

146 variable_key: str, 

147 patch_body: VariableBody, 

148 session: SessionDep, 

149 update_mask: list[str] | None = Query(None), 

150) -> VariableResponse: 

151 """Update a variable by key.""" 

152 if patch_body.key != variable_key: 

153 raise HTTPException( 

154 status.HTTP_400_BAD_REQUEST, "Invalid body, key from request body doesn't match uri parameter" 

155 ) 

156 old_variable = session.scalar(select(Variable).filter_by(key=variable_key).limit(1)) 

157 if not old_variable: 157 ↛ 158line 157 didn't jump to line 158 because the condition on line 157 was never true

158 raise HTTPException( 

159 status.HTTP_404_NOT_FOUND, f"The Variable with key: `{variable_key}` was not found" 

160 ) 

161 variable = update_orm_from_pydantic(old_variable, patch_body, update_mask) 

162 return variable 

163 

164 

165@variables_router.post( 

166 "", 

167 status_code=status.HTTP_201_CREATED, 

168 responses=create_openapi_http_exception_doc([status.HTTP_409_CONFLICT]), 

169 dependencies=[Depends(action_logging()), Depends(requires_access_variable("POST"))], 

170) 

171def post_variable( 

172 post_body: VariableBody, 

173 session: SessionDep, 

174) -> VariableResponse: 

175 """Create a variable.""" 

176 # Check if the key already exists 

177 existing_variable = session.scalar(select(Variable).where(Variable.key == post_body.key).limit(1)) 

178 if existing_variable: 

179 raise HTTPException( 

180 status_code=status.HTTP_409_CONFLICT, 

181 detail=f"The Variable with key: `{post_body.key}` already exists", 

182 ) 

183 

184 Variable.set(**post_body.model_dump(), session=session) 

185 

186 return session.scalars(select(Variable).where(Variable.key == post_body.key)).one() 

187 

188 

189@variables_router.patch( 

190 "", dependencies=[Depends(action_logging()), Depends(requires_access_variable_bulk())] 

191) 

192def bulk_variables( 

193 request: BulkBody[VariableBody], 

194 session: SessionDep, 

195) -> BulkResponse: 

196 """Bulk create, update, and delete variables.""" 

197 return BulkVariableService(session=session, request=request).handle_request()