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
« 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, Query, status
22from sqlalchemy import delete, select
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
52variables_router = AirflowRouter(tags=["Variable"], prefix="/variables")
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 )
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))
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 )
93 return variable
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 )
127 variables = session.scalars(variable_select)
129 return VariableCollectionResponse(
130 variables=variables,
131 total_entries=total_entries,
132 )
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
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 )
184 Variable.set(**post_body.model_dump(), session=session)
186 return session.scalars(select(Variable).where(Variable.key == post_body.key)).one()
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()