Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/auth/managers/simple/simple_auth_manager.py: 62%

209 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 14:22 +0000

1# 

2# Licensed to the Apache Software Foundation (ASF) under one 

3# or more contributor license agreements. See the NOTICE file 

4# distributed with this work for additional information 

5# regarding copyright ownership. The ASF licenses this file 

6# to you under the Apache License, Version 2.0 (the 

7# "License"); you may not use this file except in compliance 

8# with the License. You may obtain a copy of the License at 

9# 

10# http://www.apache.org/licenses/LICENSE-2.0 

11# 

12# Unless required by applicable law or agreed to in writing, 

13# software distributed under the License is distributed on an 

14# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY 

15# KIND, either express or implied. See the License for the 

16# specific language governing permissions and limitations 

17# under the License. 

18from __future__ import annotations 

19 

20import fcntl 

21import json 

22import logging 

23import os 

24import secrets 

25from collections import namedtuple 

26from enum import Enum 

27from json import JSONDecodeError 

28from pathlib import Path 

29from typing import TYPE_CHECKING, Any, TextIO 

30from urllib.parse import urlencode 

31 

32from fastapi import FastAPI, Request 

33from fastapi.responses import HTMLResponse 

34from fastapi.staticfiles import StaticFiles 

35from fastapi.templating import Jinja2Templates 

36from termcolor import colored 

37 

38from airflow.api_fastapi.app import AUTH_MANAGER_FASTAPI_APP_PREFIX 

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

40from airflow.api_fastapi.auth.managers.models.resource_details import TeamDetails 

41from airflow.api_fastapi.auth.managers.simple.user import SimpleAuthManagerUser 

42from airflow.api_fastapi.common.types import MenuItem 

43from airflow.configuration import AIRFLOW_HOME, conf 

44 

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

46 from starlette.middleware import _MiddlewareFactory 

47 

48 from airflow.api_fastapi.auth.managers.base_auth_manager import ResourceMethod 

49 from airflow.api_fastapi.auth.managers.models.resource_details import ( 

50 AccessView, 

51 AssetAliasDetails, 

52 AssetDetails, 

53 ConfigurationDetails, 

54 ConnectionDetails, 

55 DagAccessEntity, 

56 DagDetails, 

57 PoolDetails, 

58 VariableDetails, 

59 ) 

60 

61log = logging.getLogger(__name__) 

62 

63 

64class SimpleAuthManagerRole(namedtuple("SimpleAuthManagerRole", "name order"), Enum): 

65 """ 

66 List of pre-defined roles in simple auth manager. 

67 

68 The first attribute defines the name that references this role in the config. 

69 The second attribute defines the order between roles. The role with order X means it grants access to 

70 resources under its umbrella and all resources under the umbrella of roles of lower order 

71 """ 

72 

73 # VIEWER role gives all read-only permissions 

74 VIEWER = "VIEWER", 0 

75 

76 # USER role gives viewer role permissions + access to Dags 

77 USER = "USER", 1 

78 

79 # OP role gives user role permissions + access to connections, config, pools, variables 

80 OP = "OP", 2 

81 

82 # ADMIN role gives all permissions 

83 ADMIN = "ADMIN", 3 

84 

85 

86class SimpleAuthManager(BaseAuthManager[SimpleAuthManagerUser]): 

87 """ 

88 Simple auth manager. 

89 

90 Default auth manager used in Airflow. This auth manager should not be used in production. 

91 This auth manager is very basic and only intended for development and testing purposes. 

92 """ 

93 

94 @staticmethod 

95 def get_generated_password_file() -> str: 

96 if configured_file := conf.get("core", "simple_auth_manager_passwords_file", fallback=None): 96 ↛ 99line 96 didn't jump to line 99 because the condition on line 96 was always true

97 return configured_file 

98 

99 return os.path.join(AIRFLOW_HOME, "simple_auth_manager_passwords.json.generated") 

100 

101 @staticmethod 

102 def get_users() -> list[SimpleAuthManagerUser]: 

103 config_users = [u.split(":") for u in conf.getlist("core", "simple_auth_manager_users")] 

104 users = [] 

105 for user in config_users: 

106 teams = None 

107 if len(user) == 3: 107 ↛ 108line 107 didn't jump to line 108 because the condition on line 107 was never true

108 if not conf.getboolean("core", "multi_team"): 

109 raise ValueError( 

110 f"The user '{user[0]}' is associated to at least one team and multi-team mode is not configured in the Airflow environment." 

111 ) 

112 teams = user[2].split("|") 

113 

114 users.append(SimpleAuthManagerUser(username=user[0], role=user[1], teams=teams)) 

115 return users 

116 

117 @staticmethod 

118 def get_passwords() -> dict[str, str]: 

119 password_file = SimpleAuthManager.get_generated_password_file() 

120 with open(password_file, "r+") as file: 

121 return SimpleAuthManager._get_passwords(file) 

122 

123 @staticmethod 

124 def _looks_like_production( 

125 *, 

126 sql_conn: str | None = None, 

127 api_host: str | None = None, 

128 executor: str | None = None, 

129 ) -> bool: 

130 """ 

131 Best-effort heuristic for whether the Airflow deployment looks production-shaped. 

132 

133 Returns True if any of the following hold: 

134 

135 - The SQL backend is not sqlite (i.e. Postgres or MySQL is configured). 

136 - The API host is bound to a non-local address. 

137 - The configured executor is not Local-/Sequential-/Debug-/InProcessExecutor. 

138 

139 None of these are *definitive* — a developer can pick any combination locally 

140 — but the cumulative signal is strong enough to justify a loud warning that 

141 SimpleAuthManager (which is dev-only by design) is being used in a setup 

142 that resembles production. 

143 

144 Each axis can be passed in directly (kwargs) so unit tests can probe the 

145 decision logic without touching the global ``conf`` state. ``None`` (the 

146 default) reads the value from ``conf``. 

147 """ 

148 if sql_conn is None: 148 ↛ 150line 148 didn't jump to line 150 because the condition on line 148 was always true

149 sql_conn = conf.get("database", "sql_alchemy_conn", fallback="") 

150 if api_host is None: 150 ↛ 152line 150 didn't jump to line 152 because the condition on line 150 was always true

151 api_host = conf.get("api", "host", fallback="localhost") 

152 if executor is None: 152 ↛ 155line 152 didn't jump to line 155 because the condition on line 152 was always true

153 executor = conf.get("core", "executor", fallback="LocalExecutor") 

154 

155 if sql_conn and not sql_conn.startswith("sqlite:"): 155 ↛ 157line 155 didn't jump to line 157 because the condition on line 155 was always true

156 return True 

157 if api_host.strip() not in {"localhost", "127.0.0.1", "::1", "[::1]"}: 

158 return True 

159 # Split on '.' to get the class name only (handles fully-qualified executor paths). 

160 local_executors = {"LocalExecutor", "SequentialExecutor", "DebugExecutor", "InProcessExecutor"} 

161 if executor.split(".")[-1] not in local_executors: 

162 return True 

163 return False 

164 

165 def init(self) -> None: 

166 super().init() 

167 is_simple_auth_manager_all_admins = conf.getboolean("core", "simple_auth_manager_all_admins") 

168 if is_simple_auth_manager_all_admins: 168 ↛ 169line 168 didn't jump to line 169 because the condition on line 168 was never true

169 return 

170 

171 # SimpleAuthManager is dev-only by design — it stores passwords in plaintext, 

172 # prints generated passwords to stdout/logs on first init, and provides no 

173 # rotation mechanism. Emit a loud warning when the deployment shape suggests 

174 # production so it shows up in startup logs of misconfigured deployments. 

175 if self._looks_like_production(): 175 ↛ 185line 175 didn't jump to line 185 because the condition on line 175 was always true

176 log.warning( 

177 "SimpleAuthManager is active but the deployment shape looks like production " 

178 "(non-sqlite backend, non-local API host, or a distributed executor). " 

179 "SimpleAuthManager stores passwords in plaintext at %s and prints generated " 

180 "passwords to stdout/logs on first init. Use a real auth manager " 

181 "(e.g. FAB or Keycloak) for production deployments.", 

182 self.get_generated_password_file(), 

183 ) 

184 

185 users = self.get_users() 

186 password_file = self.get_generated_password_file() 

187 

188 try: 

189 with open(password_file, "a+") as file: 

190 try: 

191 # Non-blocking exclusive lock on this file 

192 # Fastapi spins up N workers, so this method is called N times in N different processes 

193 # This needs to be called only once so we use the file ``password_file`` as locking mechanism 

194 fcntl.flock(file, fcntl.LOCK_EX | fcntl.LOCK_NB) 

195 passwords = self._get_passwords(stream=file) 

196 changed = False 

197 for user in users: 

198 if user.username not in passwords: 198 ↛ 200line 198 didn't jump to line 200 because the condition on line 198 was never true

199 # User does not exist in the file, adding it 

200 passwords[user.username] = self._generate_password() 

201 self._print_output( 

202 f"Password for user '{user.username}': {passwords[user.username]}" 

203 ) 

204 changed = True 

205 

206 if changed: 206 ↛ 207line 206 didn't jump to line 207 because the condition on line 206 was never true

207 file.seek(0) 

208 file.truncate() 

209 file.write(json.dumps(passwords) + "\n") 

210 finally: 

211 # Release lock 

212 fcntl.flock(file, fcntl.LOCK_UN) 

213 except BlockingIOError: 

214 # The file is locked, another process called this method already, skipping 

215 pass 

216 

217 def get_url_login(self, **kwargs) -> str: 

218 """Return the login page url.""" 

219 next_url = kwargs.get("next_url") 

220 is_simple_auth_manager_all_admins = conf.getboolean("core", "simple_auth_manager_all_admins") 

221 if is_simple_auth_manager_all_admins: 221 ↛ 222line 221 didn't jump to line 222 because the condition on line 221 was never true

222 login_url = AUTH_MANAGER_FASTAPI_APP_PREFIX + "/token/login" 

223 else: 

224 login_url = AUTH_MANAGER_FASTAPI_APP_PREFIX + "/login" 

225 

226 if next_url: 226 ↛ 227line 226 didn't jump to line 227 because the condition on line 226 was never true

227 return f"{login_url}?{urlencode({'next': next_url})}" 

228 

229 return login_url 

230 

231 def deserialize_user(self, token: dict[str, Any]) -> SimpleAuthManagerUser: 

232 return SimpleAuthManagerUser( 

233 username=token["sub"], 

234 role=token["role"], 

235 teams=token.get("teams"), 

236 ) 

237 

238 def serialize_user(self, user: SimpleAuthManagerUser) -> dict[str, Any]: 

239 return {"sub": user.username, "role": user.role, "teams": user.teams} 

240 

241 def is_authorized_configuration( 

242 self, 

243 *, 

244 method: ResourceMethod, 

245 user: SimpleAuthManagerUser, 

246 details: ConfigurationDetails | None = None, 

247 ) -> bool: 

248 return self._is_authorized( 

249 method=method, 

250 allow_get_role=SimpleAuthManagerRole.VIEWER, 

251 allow_role=SimpleAuthManagerRole.OP, 

252 user=user, 

253 ) 

254 

255 def is_authorized_connection( 

256 self, 

257 *, 

258 method: ResourceMethod, 

259 user: SimpleAuthManagerUser, 

260 details: ConnectionDetails | None = None, 

261 ) -> bool: 

262 return self._is_authorized( 

263 method=method, 

264 allow_role=SimpleAuthManagerRole.OP, 

265 user=user, 

266 team_name=details.team_name if details else None, 

267 ) 

268 

269 def is_authorized_dag( 

270 self, 

271 *, 

272 method: ResourceMethod, 

273 user: SimpleAuthManagerUser, 

274 access_entity: DagAccessEntity | None = None, 

275 details: DagDetails | None = None, 

276 ) -> bool: 

277 return self._is_authorized( 

278 method=method, 

279 allow_get_role=SimpleAuthManagerRole.VIEWER, 

280 allow_role=SimpleAuthManagerRole.USER, 

281 user=user, 

282 team_name=details.team_name if details else None, 

283 ) 

284 

285 def is_authorized_asset( 

286 self, 

287 *, 

288 method: ResourceMethod, 

289 user: SimpleAuthManagerUser, 

290 details: AssetDetails | None = None, 

291 ) -> bool: 

292 return self._is_authorized( 

293 method=method, 

294 allow_get_role=SimpleAuthManagerRole.VIEWER, 

295 allow_role=SimpleAuthManagerRole.OP, 

296 user=user, 

297 ) 

298 

299 def is_authorized_asset_alias( 

300 self, 

301 *, 

302 method: ResourceMethod, 

303 user: SimpleAuthManagerUser, 

304 details: AssetAliasDetails | None = None, 

305 ) -> bool: 

306 return self._is_authorized( 

307 method=method, 

308 allow_get_role=SimpleAuthManagerRole.VIEWER, 

309 allow_role=SimpleAuthManagerRole.OP, 

310 user=user, 

311 ) 

312 

313 def is_authorized_pool( 

314 self, 

315 *, 

316 method: ResourceMethod, 

317 user: SimpleAuthManagerUser, 

318 details: PoolDetails | None = None, 

319 ) -> bool: 

320 return self._is_authorized( 

321 method=method, 

322 allow_get_role=SimpleAuthManagerRole.VIEWER, 

323 allow_role=SimpleAuthManagerRole.OP, 

324 user=user, 

325 team_name=details.team_name if details else None, 

326 ) 

327 

328 def is_authorized_team( 

329 self, 

330 *, 

331 method: ResourceMethod, 

332 user: SimpleAuthManagerUser, 

333 details: TeamDetails | None = None, 

334 ) -> bool: 

335 if not details: 

336 return False 

337 if self._is_admin(user): 

338 return True 

339 return details.name in user.teams 

340 

341 def is_authorized_variable( 

342 self, 

343 *, 

344 method: ResourceMethod, 

345 user: SimpleAuthManagerUser, 

346 details: VariableDetails | None = None, 

347 ) -> bool: 

348 return self._is_authorized( 

349 method=method, 

350 allow_role=SimpleAuthManagerRole.OP, 

351 user=user, 

352 team_name=details.team_name if details else None, 

353 ) 

354 

355 def is_authorized_view(self, *, access_view: AccessView, user: SimpleAuthManagerUser) -> bool: 

356 return self._is_authorized(method="GET", allow_role=SimpleAuthManagerRole.VIEWER, user=user) 

357 

358 def is_authorized_custom_view( 

359 self, *, method: ResourceMethod, resource_name: str, user: SimpleAuthManagerUser 

360 ): 

361 return self._is_authorized(method="GET", allow_role=SimpleAuthManagerRole.VIEWER, user=user) 

362 

363 def filter_authorized_menu_items( 

364 self, menu_items: list[MenuItem], *, user: SimpleAuthManagerUser 

365 ) -> list[MenuItem]: 

366 return menu_items 

367 

368 def is_authorized_hitl_task(self, *, assigned_users: set[str], user: SimpleAuthManagerUser) -> bool: 

369 """ 

370 Check if a user is allowed to approve/reject a HITL task. 

371 

372 When simple_auth_manager_all_admins=True, all authenticated users are allowed 

373 to approve/reject any task. Otherwise, the user must be in the assigned_users set. 

374 """ 

375 is_simple_auth_manager_all_admins = conf.getboolean("core", "simple_auth_manager_all_admins") 

376 

377 if is_simple_auth_manager_all_admins: 

378 # In all-admin mode, everyone is allowed 

379 return True 

380 

381 # Delegate to parent class for the actual authorization check 

382 return super().is_authorized_hitl_task(assigned_users=assigned_users, user=user) 

383 

384 def get_fastapi_middlewares(self) -> list[tuple[_MiddlewareFactory[Any], dict[str, Any]]]: 

385 """Register the all-admins middleware when ``[core] simple_auth_manager_all_admins`` is set.""" 

386 if not conf.getboolean("core", "simple_auth_manager_all_admins"): 386 ↛ 388line 386 didn't jump to line 388 because the condition on line 386 was always true

387 return [] 

388 from airflow.api_fastapi.auth.managers.simple.middleware import SimpleAllAdminMiddleware 

389 

390 return [(SimpleAllAdminMiddleware, {})] 

391 

392 def get_fastapi_app(self) -> FastAPI | None: 

393 """ 

394 Specify a sub FastAPI application specific to the auth manager. 

395 

396 This sub application, if specified, is mounted in the main FastAPI application. 

397 """ 

398 from airflow.api_fastapi.auth.managers.simple.routes.login import login_router 

399 

400 dev_mode = os.environ.get("DEV_MODE", str(False)) == "true" 

401 directory = Path(__file__).parent.joinpath("ui", "dev" if dev_mode else "dist") 

402 directory.mkdir(exist_ok=True) 

403 

404 templates = Jinja2Templates(directory=directory) 

405 

406 app = FastAPI( 

407 title="Simple auth manager sub application", 

408 description=( 

409 "This is the simple auth manager fastapi sub application. This API is only available if the " 

410 "auth manager used in the Airflow environment is simple auth manager. " 

411 "This sub application provides the login form for users to log in." 

412 ), 

413 ) 

414 app.include_router(login_router) 

415 app.mount( 

416 "/static", 

417 StaticFiles( 

418 directory=directory, 

419 html=True, 

420 ), 

421 name="simple_auth_manager_ui_folder", 

422 ) 

423 

424 @app.get("/{rest_of_path:path}", response_class=HTMLResponse, include_in_schema=False) 

425 def webapp(request: Request, rest_of_path: str): 

426 return templates.TemplateResponse( 

427 request, 

428 "/index.html", 

429 {"backend_server_base_url": request.base_url.path}, 

430 media_type="text/html", 

431 ) 

432 

433 return app 

434 

435 def _get_teams(self) -> set[str]: 

436 users = self.get_users() 

437 return {team for user in users for team in user.teams} 

438 

439 @staticmethod 

440 def _is_admin(user: SimpleAuthManagerUser) -> bool: 

441 """Return whether the user has the Admin role.""" 

442 if not user.role: 442 ↛ 443line 442 didn't jump to line 443 because the condition on line 442 was never true

443 return False 

444 

445 role_str = user.role.upper() 

446 role = SimpleAuthManagerRole[role_str] 

447 

448 return role == SimpleAuthManagerRole.ADMIN 

449 

450 @staticmethod 

451 def _is_authorized( 

452 *, 

453 method: ResourceMethod, 

454 allow_role: SimpleAuthManagerRole, 

455 user: SimpleAuthManagerUser, 

456 allow_get_role: SimpleAuthManagerRole | None = None, 

457 team_name: str | None = None, 

458 ): 

459 """ 

460 Return whether the user is authorized to access a given resource. 

461 

462 :param method: the method to perform 

463 :param allow_role: minimal role giving access to the resource, if the user's role is greater or 

464 equal than this role, they have access 

465 :param user: the user to check the authorization for 

466 :param allow_get_role: minimal role giving access to the resource, if the user's role is greater or 

467 equal than this role, they have access. If not provided, ``allow_role`` is used 

468 :param team_name: team associated to the resource (if any) 

469 """ 

470 if not user.role: 470 ↛ 471line 470 didn't jump to line 471 because the condition on line 470 was never true

471 return False 

472 

473 if SimpleAuthManager._is_admin(user): 473 ↛ 476line 473 didn't jump to line 476 because the condition on line 473 was always true

474 return True 

475 

476 if team_name and team_name not in user.teams: 

477 return False 

478 

479 if not allow_get_role: 

480 allow_get_role = allow_role 

481 

482 role_str = user.role.upper() 

483 role = SimpleAuthManagerRole[role_str] 

484 if method == "GET": 

485 return role.order >= allow_get_role.order 

486 return role.order >= allow_role.order 

487 

488 @staticmethod 

489 def _get_passwords(stream: TextIO) -> dict[str, str]: 

490 try: 

491 # Read passwords from file 

492 stream.seek(0) 

493 content = stream.read().strip() or "{}" 

494 user_passwords_from_file = json.loads(content) 

495 except JSONDecodeError: 

496 log.error("Error decoding JSON from file %s", stream.name) 

497 raise 

498 

499 return user_passwords_from_file 

500 

501 @staticmethod 

502 def _generate_password() -> str: 

503 alphabet = "abcdefghkmnpqrstuvwxyzABCDEFGHKMNPQRSTUVWXYZ23456789" 

504 return "".join(secrets.choice(alphabet) for _ in range(16)) 

505 

506 @staticmethod 

507 def _print_output(output: str): 

508 if conf.getboolean("logging", "json_logs", fallback=False): 

509 for line in output.splitlines(): 

510 log.info(line.strip()) 

511 else: 

512 name = "Simple auth manager" 

513 colorized_name = colored(f"{name:10}", "white") 

514 for line in output.splitlines(): 

515 print(f"{colorized_name} | {line.strip()}")