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
« 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
19import logging
20import threading
21from contextlib import AsyncExitStack, asynccontextmanager
22from functools import cache
23from typing import TYPE_CHECKING
24from urllib.parse import urlsplit
26from fastapi import FastAPI
27from fastapi.routing import Mount
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
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
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
51# Define the full path on which the potential auth manager fastapi is mounted
52AUTH_MANAGER_FASTAPI_APP_PREFIX = f"{API_ROOT_PATH}auth"
55def get_cookie_path() -> str:
56 """
57 Return the path to scope cookies to, derived from ``[api] base_url``.
59 Falls back to ``"/"`` when no ``base_url`` is configured.
60 """
61 return API_ROOT_PATH or "/"
64# Fast API apps mounted under these prefixes are not allowed
65RESERVED_URL_PREFIXES = ["/api/v2", "/ui", "/execution", "/auth", "/pluginsv2"]
67log = logging.getLogger(__name__)
70class _AuthManagerState:
71 instance: BaseAuthManager | None = None
72 _lock = threading.Lock()
75def _initialize_api_server_stats() -> None:
76 """
77 Initialize the ``Stats`` singleton in the API server process.
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
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 )
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
111@providers_configuration_loaded
112def create_app(apps: str = "all") -> FastAPI:
113 apps_list = apps.split(",") if apps else ["all"]
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 )
139 dag_bag = create_dag_bag()
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)
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)
155 init_access_logging(app)
157 init_config(app)
159 return app
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)
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
174def get_auth_manager_cls() -> type[BaseAuthManager]:
175 """
176 Return just the auth manager class without initializing it.
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")
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 )
187 return auth_manager_cls
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
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
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)
211 return am
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
224def init_access_logging(app: FastAPI) -> None:
225 """Install the access log middleware, the only producer of access records."""
226 app.add_middleware(HttpAccessLogMiddleware)
229def init_plugins(app: FastAPI) -> None:
230 """Integrate FastAPI app, middlewares and UI plugins."""
231 from airflow import plugins_manager
233 apps, root_middlewares = plugins_manager.get_fastapi_plugins()
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
252 log.debug("Adding subapplication %s under prefix %s", name, url_prefix)
253 app.mount(url_prefix, subapp)
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", {})
261 if middleware is None:
262 log.error("'middleware' key is missing for the fastapi middleware: %s", name)
263 continue
265 if not callable(middleware):
266 log.error("'middleware' value for %s is should be callable: %s", name, middleware)
267 continue
269 log.debug("Adding root middleware %s", name)
270 app.add_middleware(middleware, *args, **kwargs)