Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/app.py: 64%

134 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 

19import logging 

20import threading 

21from contextlib import AsyncExitStack, asynccontextmanager 

22from functools import cache 

23from typing import TYPE_CHECKING 

24from urllib.parse import urlsplit 

25 

26from fastapi import FastAPI 

27from fastapi.routing import Mount 

28 

29from airflow.api_fastapi.common.dagbag import create_dag_bag 

30from airflow.api_fastapi.common.exceptions import init_error_handlers 

31from airflow.api_fastapi.common.http_access_log import HttpAccessLogMiddleware 

32from airflow.api_fastapi.core_api.app import ( 

33 init_config, 

34 init_flask_plugins, 

35 init_middlewares, 

36 init_views, 

37) 

38from airflow.api_fastapi.execution_api.app import create_task_execution_api_app 

39from airflow.configuration import conf 

40from airflow.exceptions import AirflowConfigException 

41from airflow.utils.providers_configuration_loader import providers_configuration_loaded 

42 

43if TYPE_CHECKING: 43 ↛ 44line 43 didn't jump to line 44 because the condition on line 43 was never true

44 from airflow.api_fastapi.auth.managers.base_auth_manager import BaseAuthManager 

45 

46API_BASE_URL = conf.get("api", "base_url", fallback="") 

47if not API_BASE_URL or not API_BASE_URL.endswith("/"): 47 ↛ 49line 47 didn't jump to line 49 because the condition on line 47 was always true

48 API_BASE_URL += "/" 

49API_ROOT_PATH = urlsplit(API_BASE_URL).path 

50 

51# Define the full path on which the potential auth manager fastapi is mounted 

52AUTH_MANAGER_FASTAPI_APP_PREFIX = f"{API_ROOT_PATH}auth" 

53 

54 

55def get_cookie_path() -> str: 

56 """ 

57 Return the path to scope cookies to, derived from ``[api] base_url``. 

58 

59 Falls back to ``"/"`` when no ``base_url`` is configured. 

60 """ 

61 return API_ROOT_PATH or "/" 

62 

63 

64# Fast API apps mounted under these prefixes are not allowed 

65RESERVED_URL_PREFIXES = ["/api/v2", "/ui", "/execution", "/auth", "/pluginsv2"] 

66 

67log = logging.getLogger(__name__) 

68 

69 

70class _AuthManagerState: 

71 instance: BaseAuthManager | None = None 

72 _lock = threading.Lock() 

73 

74 

75def _initialize_api_server_stats() -> None: 

76 """ 

77 Initialize the ``Stats`` singleton in the API server process. 

78 

79 Initialization is guarded so a metrics misconfiguration can never prevent the API server 

80 from starting. 

81 """ 

82 try: 

83 from airflow._shared.observability.metrics import stats 

84 from airflow.observability.metrics import stats_utils 

85 

86 stats.initialize( 

87 factory=stats_utils.get_stats_factory(), 

88 export_legacy_names=conf.getboolean("metrics", "legacy_names_on"), 

89 ) 

90 except Exception: 

91 log.warning( 

92 "Failed to initialize API server Stats in the API server; metrics emitted through the " 

93 "API server Stats singleton will not be recorded.", 

94 exc_info=True, 

95 ) 

96 

97 

98@asynccontextmanager 

99async def lifespan(app: FastAPI): 

100 _initialize_api_server_stats() 

101 async with AsyncExitStack() as stack: 

102 for route in app.routes: 

103 if isinstance(route, Mount) and isinstance(route.app, FastAPI): 

104 await stack.enter_async_context( 

105 route.app.router.lifespan_context(route.app), 

106 ) 

107 app.state.lifespan_called = True 

108 yield 

109 

110 

111@providers_configuration_loaded 

112def create_app(apps: str = "all") -> FastAPI: 

113 apps_list = apps.split(",") if apps else ["all"] 

114 

115 app = FastAPI( 

116 title="Airflow API", 

117 description="Airflow API. All endpoints located under ``/api/v2`` can be used safely, are stable and backward compatible. " 

118 "Endpoints located under ``/ui`` are dedicated to the UI and are subject to breaking change " 

119 "depending on the need of the frontend. Users should not rely on those but use the public ones instead." 

120 "\n\n" 

121 "**Filtering with pattern parameters.** Many list endpoints accept ``*_pattern`` and " 

122 "``*_prefix_pattern`` query parameters. Unless a parameter's own description says otherwise, " 

123 "``*_pattern`` is a case-insensitive substring match (SQL ``ILIKE '%term%'``) where ``%`` matches " 

124 "any sequence and ``_`` matches any single character (e.g. ``%customer_%``) — convenient, but it " 

125 "cannot use B-tree indexes, so it is slow on large tables. ``*_prefix_pattern`` matches the start " 

126 "of the value, is case-sensitive and index-friendly (prefer it at scale); there ``%`` and ``_`` " 

127 "are literal and trailing non-alphanumeric characters are stripped so the range scan stays " 

128 "index-compatible under locale-aware collations " 

129 "(e.g. ``test_`` matches values starting with ``test``, and ``s3://`` matches ``s3``). In both, " 

130 "``|`` means OR (e.g. ``dag1|dag2``) and ``~`` matches everything. Regular expressions are not " 

131 "supported by these parameters; regex-capable endpoints expose a separate parameter.", 

132 lifespan=lifespan, 

133 root_path=API_ROOT_PATH.removesuffix("/"), 

134 version="2", 

135 docs_url="/docs" if conf.getboolean("api", "enable_swagger_ui") else None, 

136 redoc_url="/redoc" if conf.getboolean("api", "enable_swagger_ui") else None, 

137 ) 

138 

139 dag_bag = create_dag_bag() 

140 

141 if "all" in apps_list or "execution" in apps_list: 141 ↛ 146line 141 didn't jump to line 146 because the condition on line 141 was always true

142 task_exec_api_app = create_task_execution_api_app() 

143 task_exec_api_app.state.dag_bag = dag_bag 

144 app.mount("/execution", task_exec_api_app) 

145 

146 if "all" in apps_list or "core" in apps_list: 146 ↛ 155line 146 didn't jump to line 155 because the condition on line 146 was always true

147 app.state.dag_bag = dag_bag 

148 init_plugins(app) 

149 init_auth_manager(app) 

150 init_flask_plugins(app) 

151 init_views(app) # Core views need to be the last routes added - it has a catch all route 

152 init_error_handlers(app) 

153 init_middlewares(app) 

154 

155 init_access_logging(app) 

156 

157 init_config(app) 

158 

159 return app 

160 

161 

162@cache 

163def cached_app(config=None, testing=False, apps="all") -> FastAPI: 

164 """Return cached instance of Airflow API app.""" 

165 return create_app(apps=apps) 

166 

167 

168def purge_cached_app() -> None: 

169 """Remove the cached version of the app and auth_manager in global state.""" 

170 cached_app.cache_clear() 

171 _AuthManagerState.instance = None 

172 

173 

174def get_auth_manager_cls() -> type[BaseAuthManager]: 

175 """ 

176 Return just the auth manager class without initializing it. 

177 

178 Useful to save execution time if only static methods need to be called. 

179 """ 

180 auth_manager_cls = conf.getimport(section="core", key="auth_manager") 

181 

182 if not auth_manager_cls: 182 ↛ 183line 182 didn't jump to line 183 because the condition on line 182 was never true

183 raise AirflowConfigException( 

184 "No auth manager defined in the config. Please specify one using section/key [core/auth_manager]." 

185 ) 

186 

187 return auth_manager_cls 

188 

189 

190def create_auth_manager() -> BaseAuthManager: 

191 """Create the auth manager, cached as a thread-safe singleton.""" 

192 auth_manager_cls = get_auth_manager_cls() 

193 if _AuthManagerState.instance is not None and isinstance(_AuthManagerState.instance, auth_manager_cls): 193 ↛ 194line 193 didn't jump to line 194 because the condition on line 193 was never true

194 return _AuthManagerState.instance 

195 with _AuthManagerState._lock: 

196 if _AuthManagerState.instance is None or not isinstance(_AuthManagerState.instance, auth_manager_cls): 196 ↛ 198line 196 didn't jump to line 198

197 _AuthManagerState.instance = auth_manager_cls() 

198 return _AuthManagerState.instance 

199 

200 

201def init_auth_manager(app: FastAPI | None = None) -> BaseAuthManager: 

202 """Initialize the auth manager.""" 

203 am = create_auth_manager() 

204 am.init() 

205 if app: 205 ↛ 208line 205 didn't jump to line 208 because the condition on line 205 was always true

206 app.state.auth_manager = am 

207 

208 if app and (auth_manager_fastapi_app := am.get_fastapi_app()): 208 ↛ 211line 208 didn't jump to line 211 because the condition on line 208 was always true

209 app.mount("/auth", auth_manager_fastapi_app) 

210 

211 return am 

212 

213 

214def get_auth_manager() -> BaseAuthManager: 

215 """Return the auth manager, provided it's been initialized before.""" 

216 if _AuthManagerState.instance is None: 216 ↛ 217line 216 didn't jump to line 217 because the condition on line 216 was never true

217 raise RuntimeError( 

218 "Auth Manager has not been initialized yet. " 

219 "The `init_auth_manager` method needs to be called first." 

220 ) 

221 return _AuthManagerState.instance 

222 

223 

224def init_access_logging(app: FastAPI) -> None: 

225 """Install the access log middleware, the only producer of access records.""" 

226 app.add_middleware(HttpAccessLogMiddleware) 

227 

228 

229def init_plugins(app: FastAPI) -> None: 

230 """Integrate FastAPI app, middlewares and UI plugins.""" 

231 from airflow import plugins_manager 

232 

233 apps, root_middlewares = plugins_manager.get_fastapi_plugins() 

234 

235 for subapp_dict in apps: 235 ↛ 236line 235 didn't jump to line 236 because the loop on line 235 never started

236 name = subapp_dict.get("name") 

237 subapp = subapp_dict.get("app") 

238 if subapp is None: 

239 log.error("'app' key is missing for the fastapi app: %s", name) 

240 continue 

241 url_prefix = subapp_dict.get("url_prefix") 

242 if url_prefix is None: 

243 log.error("'url_prefix' key is missing for the fastapi app: %s", name) 

244 continue 

245 if url_prefix == "": 

246 log.error("'url_prefix' key is empty string for the fastapi app: %s", name) 

247 continue 

248 if any(url_prefix.startswith(prefix) for prefix in RESERVED_URL_PREFIXES): 

249 log.error("Plugin %s attempted to use reserved url_prefix '%s'", name, url_prefix) 

250 continue 

251 

252 log.debug("Adding subapplication %s under prefix %s", name, url_prefix) 

253 app.mount(url_prefix, subapp) 

254 

255 for middleware_dict in root_middlewares: 255 ↛ 256line 255 didn't jump to line 256 because the loop on line 255 never started

256 name = middleware_dict.get("name") 

257 middleware = middleware_dict.get("middleware") 

258 args = middleware_dict.get("args", []) 

259 kwargs = middleware_dict.get("kwargs", {}) 

260 

261 if middleware is None: 

262 log.error("'middleware' key is missing for the fastapi middleware: %s", name) 

263 continue 

264 

265 if not callable(middleware): 

266 log.error("'middleware' value for %s is should be callable: %s", name, middleware) 

267 continue 

268 

269 log.debug("Adding root middleware %s", name) 

270 app.add_middleware(middleware, *args, **kwargs)