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
« 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.
18from __future__ import annotations
20from typing import cast
22from fastapi import HTTPException, status
23from fastapi.exceptions import RequestValidationError
24from pydantic import ValidationError
25from sqlalchemy import select
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
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.
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"}
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 )
80class BulkVariableService(BulkService[VariableBody]):
81 """Service for handling bulk operations on variables."""
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
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)
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
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)
118 except HTTPException as e:
119 results.errors.append({"error": f"{e.detail}", "status_code": e.status_code})
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
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 )
143 results.success.append(updated_variable.key)
145 except HTTPException as e:
146 results.errors.append({"error": f"{e.detail}", "status_code": e.status_code})
148 except ValidationError as e:
149 results.errors.append({"error": f"{e.errors()}"})
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)
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
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)
173 except HTTPException as e:
174 results.errors.append({"error": f"{e.detail}", "status_code": e.status_code})