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

79 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. 

17 

18from __future__ import annotations 

19 

20from typing import cast 

21 

22from fastapi import HTTPException, status 

23from fastapi.exceptions import RequestValidationError 

24from pydantic import ValidationError 

25from sqlalchemy import select 

26 

27from airflow.api_fastapi.core_api.datamodels.common import ( 

28 BulkActionNotOnExistence, 

29 BulkActionOnExistence, 

30 BulkActionResponse, 

31 BulkCreateAction, 

32 BulkDeleteAction, 

33 BulkUpdateAction, 

34) 

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

36 VariableBody, 

37 VariableBodyPartial, 

38) 

39from airflow.api_fastapi.core_api.services.public.common import BulkService 

40from airflow.models.variable import Variable 

41 

42 

43def update_orm_from_pydantic( 

44 old_variable: Variable, patch_body: VariableBody, update_mask: list[str] | None 

45) -> Variable: 

46 """ 

47 Update an existing Variable. 

48 

49 :param old_variable: The existing Variable ORM object to update. 

50 :param patch_body: The patch request body containing fields to update. 

51 :param update_mask: List of fields to update. If None, all provided fields will be updated. 

52 :return: The updated Variable object. 

53 :raises HTTPException: If attempting to update restricted fields (e.g., ``key``). 

54 """ 

55 if update_mask: 55 ↛ 56line 55 didn't jump to line 56 because the condition on line 55 was never true

56 fields_to_update = patch_body.model_fields_set & set(update_mask) 

57 try: 

58 VariableBodyPartial(**patch_body.model_dump(include=fields_to_update)) 

59 except ValidationError as e: 

60 raise RequestValidationError(errors=e.errors()) 

61 else: 

62 try: 

63 VariableBody(**patch_body.model_dump()) 

64 except ValidationError as e: 

65 raise RequestValidationError(errors=e.errors()) 

66 non_update_fields = {"key"} 

67 

68 # Apply patch via utility 

69 return cast( 

70 "Variable", 

71 BulkService.apply_patch_with_update_mask( 

72 model=old_variable, 

73 patch_body=patch_body, 

74 update_mask=update_mask, 

75 non_update_fields=non_update_fields, 

76 ), 

77 ) 

78 

79 

80class BulkVariableService(BulkService[VariableBody]): 

81 """Service for handling bulk operations on variables.""" 

82 

83 def categorize_keys(self, keys: set) -> tuple[dict, set, set]: 

84 """Categorize the given keys into matched_keys and not_found_keys based on existing keys.""" 

85 existing_variables = self.session.execute(select(Variable).filter(Variable.key.in_(keys))).scalars() 

86 existing_variables_dict = {variable.key: variable for variable in existing_variables} 

87 matched_keys = set(existing_variables_dict.keys()) 

88 not_found_keys = keys - matched_keys 

89 return existing_variables_dict, matched_keys, not_found_keys 

90 

91 def handle_bulk_create(self, action: BulkCreateAction, results: BulkActionResponse) -> None: 

92 """Bulk create variables.""" 

93 to_create_keys = {variable.key for variable in action.entities} 

94 _, matched_keys, not_found_keys = self.categorize_keys(to_create_keys) 

95 

96 try: 

97 if action.action_on_existence == BulkActionOnExistence.FAIL and matched_keys: 

98 raise HTTPException( 

99 status_code=status.HTTP_409_CONFLICT, 

100 detail=f"The variables with these keys: {matched_keys} already exist.", 

101 ) 

102 if action.action_on_existence == BulkActionOnExistence.SKIP: 

103 create_keys = not_found_keys 

104 else: 

105 create_keys = to_create_keys 

106 

107 for variable in action.entities: 

108 if variable.key in create_keys: 

109 # VariableBody already JSON-encodes non-string values, so no serialize_json here. 

110 Variable.set( 

111 key=variable.key, 

112 value=variable.value, 

113 description=variable.description, 

114 session=self.session, 

115 ) 

116 results.success.append(variable.key) 

117 

118 except HTTPException as e: 

119 results.errors.append({"error": f"{e.detail}", "status_code": e.status_code}) 

120 

121 def handle_bulk_update(self, action: BulkUpdateAction, results: BulkActionResponse) -> None: 

122 """Bulk Update variables.""" 

123 to_update_keys = {variable.key for variable in action.entities} 

124 existing_variables_dict, matched_keys, not_found_keys = self.categorize_keys(to_update_keys) 

125 try: 

126 if action.action_on_non_existence == BulkActionNotOnExistence.FAIL and not_found_keys: 

127 raise HTTPException( 

128 status_code=status.HTTP_404_NOT_FOUND, 

129 detail=f"The variables with these keys: {not_found_keys} were not found.", 

130 ) 

131 if action.action_on_non_existence == BulkActionNotOnExistence.SKIP: 

132 update_keys = matched_keys 

133 else: 

134 update_keys = to_update_keys 

135 

136 for variable in action.entities: 

137 if variable.key not in update_keys: 137 ↛ 138line 137 didn't jump to line 138 because the condition on line 137 was never true

138 continue 

139 updated_variable = update_orm_from_pydantic( 

140 existing_variables_dict[variable.key], variable, action.update_mask 

141 ) 

142 

143 results.success.append(updated_variable.key) 

144 

145 except HTTPException as e: 

146 results.errors.append({"error": f"{e.detail}", "status_code": e.status_code}) 

147 

148 except ValidationError as e: 

149 results.errors.append({"error": f"{e.errors()}"}) 

150 

151 def handle_bulk_delete(self, action: BulkDeleteAction, results: BulkActionResponse) -> None: 

152 """Bulk delete variables.""" 

153 to_delete_keys = set(action.entities) 

154 existing_variables_dict, matched_keys, not_found_keys = self.categorize_keys(to_delete_keys) 

155 

156 try: 

157 if action.action_on_non_existence == BulkActionNotOnExistence.FAIL and not_found_keys: 

158 raise HTTPException( 

159 status_code=status.HTTP_404_NOT_FOUND, 

160 detail=f"The variables with these keys: {not_found_keys} were not found.", 

161 ) 

162 if action.action_on_non_existence == BulkActionNotOnExistence.SKIP: 

163 delete_keys = matched_keys 

164 else: 

165 delete_keys = to_delete_keys 

166 

167 for key in delete_keys: 

168 existing_variable = existing_variables_dict.get(key) 

169 if existing_variable: 169 ↛ 167line 169 didn't jump to line 167 because the condition on line 169 was always true

170 self.session.delete(existing_variable) 

171 results.success.append(key) 

172 

173 except HTTPException as e: 

174 results.errors.append({"error": f"{e.detail}", "status_code": e.status_code})