Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/dag_sources.py: 87%

35 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 

19from fastapi import Depends, HTTPException, Response, status 

20from sqlalchemy import select 

21 

22from airflow.api_fastapi.app import get_auth_manager 

23from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity 

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

25from airflow.api_fastapi.common.headers import HeaderAcceptJsonOrText 

26from airflow.api_fastapi.common.router import AirflowRouter 

27from airflow.api_fastapi.common.types import Mimetype 

28from airflow.api_fastapi.core_api.datamodels.dag_sources import DAGSourceResponse 

29from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc 

30from airflow.api_fastapi.core_api.security import GetUserDep, requires_access_dag 

31from airflow.models import DagModel 

32from airflow.models.dag_version import DagVersion 

33 

34REDACTED_SOURCE = "REDACTED - you do not have read permission on all Dags in the file" 

35dag_sources_router = AirflowRouter(tags=["DagSource"], prefix="/dagSources") 

36 

37 

38@dag_sources_router.get( 

39 "/{dag_id}", 

40 responses={ 

41 **create_openapi_http_exception_doc( 

42 [ 

43 status.HTTP_400_BAD_REQUEST, 

44 status.HTTP_404_NOT_FOUND, 

45 status.HTTP_406_NOT_ACCEPTABLE, 

46 ] 

47 ), 

48 "200": { 

49 "description": "Successful Response", 

50 "content": { 

51 Mimetype.TEXT: {"schema": {"type": "string", "example": "dag code"}}, 

52 }, 

53 }, 

54 }, 

55 response_model=DAGSourceResponse, 

56 dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.CODE))], 

57) 

58def get_dag_source( 

59 accept: HeaderAcceptJsonOrText, 

60 dag_id: str, 

61 session: SessionDep, 

62 user: GetUserDep, 

63 version_number: int | None = None, 

64): 

65 """Get source code using file token.""" 

66 dag_version = DagVersion.get_version(dag_id, version_number, session=session) 

67 if not dag_version: 

68 raise HTTPException( 

69 status.HTTP_404_NOT_FOUND, 

70 f"The source code of the Dag {dag_id}, version_number {version_number} was not found", 

71 ) 

72 if not dag_version.dag_code: 72 ↛ 73line 72 didn't jump to line 73 because the condition on line 72 was never true

73 raise HTTPException( 

74 status.HTTP_404_NOT_FOUND, 

75 detail=f"Code not found. dag_id='{dag_id}' version_number='{version_number}'", 

76 ) 

77 

78 # Per-file authorization overlay on top of the ``DagAccessEntity.CODE`` 

79 # check above: a single source file may define multiple Dags, and the 

80 # caller having CODE access to ``dag_id`` does not imply they may read 

81 # every other Dag co-located in the file. Match the file by 

82 # ``(relative_fileloc, bundle_name)`` -- the same keying 

83 # ``import_error.py`` uses for its equivalent check -- and redact the 

84 # response when any co-located Dag is not in the caller's readable set. 

85 content = dag_version.dag_code.source_code 

86 dag_model = dag_version.dag_model 

87 if dag_model is not None and dag_model.relative_fileloc: 87 ↛ 101line 87 didn't jump to line 101 because the condition on line 87 was always true

88 file_dag_ids = set( 

89 session.scalars( 

90 select(DagModel.dag_id).where( 

91 DagModel.relative_fileloc == dag_model.relative_fileloc, 

92 DagModel.bundle_name == dag_model.bundle_name, 

93 ) 

94 ).all() 

95 ) 

96 if file_dag_ids: 96 ↛ 101line 96 didn't jump to line 101 because the condition on line 96 was always true

97 readable_dag_ids = get_auth_manager().get_authorized_dag_ids(user=user) 

98 if not file_dag_ids.issubset(readable_dag_ids): 98 ↛ 99line 98 didn't jump to line 99 because the condition on line 98 was never true

99 content = REDACTED_SOURCE 

100 

101 dag_source_model = DAGSourceResponse( 

102 dag_id=dag_id, 

103 content=content, 

104 version_number=dag_version.version_number, 

105 dag_display_name=dag_version.dag_model.dag_display_name, 

106 ) 

107 

108 if accept == Mimetype.TEXT: 

109 return Response(dag_source_model.content, media_type=Mimetype.TEXT) 

110 return dag_source_model