Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/variables.py: 63%

39 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 

20import logging 

21from typing import Annotated 

22 

23from fastapi import APIRouter, Depends, HTTPException, Path, Query, Request, status 

24from sqlalchemy import func, select 

25 

26from airflow.api_fastapi.common.db.common import SessionDep 

27from airflow.api_fastapi.execution_api.datamodels.variable import ( 

28 VariableKeysResponse, 

29 VariablePostBody, 

30 VariableResponse, 

31) 

32from airflow.api_fastapi.execution_api.security import CurrentTIToken, get_team_name_dep 

33from airflow.models.variable import Variable 

34 

35 

36async def has_variable_access( 

37 request: Request, 

38 variable_key: Annotated[str, Path(min_length=1)], 

39 token=CurrentTIToken, 

40): 

41 """Check if the task has access to the variable.""" 

42 write = request.method not in {"GET", "HEAD", "OPTIONS"} 

43 

44 log.debug( 

45 "Checking %s access for task instance with key '%s' to variable '%s'", 

46 "write" if write else "read", 

47 token.id, 

48 variable_key, 

49 ) 

50 

51 # The current version of Airflow does not support true 

52 # multi-tenancy yet (this is well-documented at 

53 # https://airflow.apache.org/docs/apache-airflow/stable/security/security_model.html#limiting-dag-author-access-to-subset-of-dags), 

54 # so for now we always return 'True' here. 

55 # When we introduce true multi-tenancy in the future 

56 # this would be the place to do add a check. 

57 return True 

58 

59 

60router = APIRouter() 

61 

62log = logging.getLogger(__name__) 

63 

64 

65# /keys must be declared before /{variable_key:path} so the static path is 

66# matched first; otherwise the catch-all path param would swallow it. 

67# has_variable_access is applied per-route below (not at router level) because 

68# it requires a variable_key path parameter that /keys does not have. 

69@router.get( 

70 "/keys", 

71 responses={ 

72 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"}, 

73 }, 

74) 

75def get_variable_keys( 

76 session: SessionDep, 

77 team_name: Annotated[str | None, Depends(get_team_name_dep)] = None, 

78 prefix: Annotated[str | None, Query()] = None, 

79 limit: Annotated[int, Query(ge=1, le=10_000)] = 1000, 

80 offset: Annotated[int, Query(ge=0)] = 0, 

81) -> VariableKeysResponse: 

82 """ 

83 Get Airflow Variable keys, optionally filtered by prefix. 

84 

85 .. note:: 

86 This endpoint deliberately bypasses the per-variable ``has_variable_access`` 

87 check, since access scoping requires a specific variable key. Any authenticated 

88 task within a team can therefore enumerate every variable key in that team — 

89 including keys for variables it would not be allowed to read. This is consistent 

90 with Airflow's security model (workers within a deployment trust each other), 

91 but the asymmetry between key enumeration and value access is intentional. 

92 """ 

93 stmt = select(Variable.key).order_by(Variable.key) 

94 if prefix is not None: 

95 stmt = stmt.where(Variable.key.startswith(prefix, autoescape=True)) 

96 if team_name is not None: 

97 stmt = stmt.where(Variable.team_name == team_name) 

98 

99 total_entries = session.scalar(select(func.count()).select_from(stmt.subquery())) or 0 

100 keys = session.scalars(stmt.offset(offset).limit(limit)).all() 

101 return VariableKeysResponse(keys=list(keys), total_entries=total_entries) 

102 

103 

104@router.get( 

105 "/{variable_key:path}", 

106 dependencies=[Depends(has_variable_access)], 

107 responses={ 

108 status.HTTP_404_NOT_FOUND: {"description": "Variable not found"}, 

109 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"}, 

110 status.HTTP_403_FORBIDDEN: {"description": "Task does not have access to the variable"}, 

111 }, 

112) 

113def get_variable( 

114 variable_key: Annotated[str, Path(min_length=1)], 

115 team_name: Annotated[str | None, Depends(get_team_name_dep)], 

116) -> VariableResponse: 

117 """Get an Airflow Variable.""" 

118 try: 

119 variable_value = Variable.get(variable_key, team_name=team_name) 

120 except KeyError: 

121 raise HTTPException( 

122 status.HTTP_404_NOT_FOUND, 

123 detail={ 

124 "reason": "not_found", 

125 "message": f"Variable with key '{variable_key}' not found", 

126 }, 

127 ) 

128 

129 return VariableResponse(key=variable_key, value=variable_value) 

130 

131 

132@router.put( 

133 "/{variable_key:path}", 

134 dependencies=[Depends(has_variable_access)], 

135 status_code=status.HTTP_201_CREATED, 

136 responses={ 

137 status.HTTP_404_NOT_FOUND: {"description": "Variable not found"}, 

138 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"}, 

139 status.HTTP_403_FORBIDDEN: {"description": "Task does not have access to the variable"}, 

140 }, 

141) 

142def put_variable( 

143 variable_key: Annotated[str, Path(min_length=1)], 

144 body: VariablePostBody, 

145 team_name: Annotated[str | None, Depends(get_team_name_dep)], 

146): 

147 """Set an Airflow Variable.""" 

148 Variable.set(key=variable_key, value=body.value, description=body.description, team_name=team_name) 

149 return {"message": "Variable successfully set"} 

150 

151 

152@router.delete( 

153 "/{variable_key:path}", 

154 dependencies=[Depends(has_variable_access)], 

155 status_code=status.HTTP_204_NO_CONTENT, 

156 responses={ 

157 status.HTTP_404_NOT_FOUND: {"description": "Variable not found"}, 

158 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"}, 

159 status.HTTP_403_FORBIDDEN: {"description": "Task does not have access to the variable"}, 

160 }, 

161) 

162def delete_variable( 

163 variable_key: Annotated[str, Path(min_length=1)], 

164 team_name: Annotated[str | None, Depends(get_team_name_dep)], 

165): 

166 """Delete an Airflow Variable.""" 

167 Variable.delete(key=variable_key, team_name=team_name)