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

84 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 contextlib 

21import textwrap 

22from collections.abc import Generator, Iterable 

23 

24from fastapi import Depends, HTTPException, Request, status 

25from fastapi.responses import StreamingResponse 

26from itsdangerous import BadSignature, URLSafeSerializer 

27from pydantic import NonNegativeInt, PositiveInt 

28from sqlalchemy.orm import joinedload 

29from sqlalchemy.sql import select 

30 

31from airflow.api_fastapi.common.dagbag import DagBagDep 

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

33from airflow.api_fastapi.common.headers import HeaderAcceptJsonOrNdjson 

34from airflow.api_fastapi.common.router import AirflowRouter 

35from airflow.api_fastapi.common.types import Mimetype 

36from airflow.api_fastapi.core_api.datamodels.log import ExternalLogUrlResponse, TaskInstancesLogResponse 

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

38from airflow.api_fastapi.core_api.security import DagAccessEntity, requires_access_dag 

39from airflow.configuration import conf 

40from airflow.exceptions import TaskNotFound 

41from airflow.models import TaskInstance, Trigger 

42from airflow.models.taskinstancehistory import TaskInstanceHistory 

43from airflow.utils.log.log_reader import TaskLogReader 

44 

45_NDJSON_BATCH_SIZE = conf.getint("api", "log_stream_buffer_size") 

46 

47task_instances_log_router = AirflowRouter( 

48 tags=["Task Instance"], prefix="/dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances" 

49) 

50 

51ndjson_example_response_for_get_log = { 

52 Mimetype.NDJSON: { 

53 "schema": { 

54 "type": "string", 

55 "example": textwrap.dedent( 

56 """\ 

57 {"content": "content"} 

58 {"content": "content"} 

59 """ 

60 ), 

61 } 

62 } 

63} 

64 

65 

66def _buffered_ndjson_stream( 

67 raw_stream: Iterable[str], 

68) -> Generator[str, None, None]: 

69 buf: list[str] = [] 

70 for line in raw_stream: 

71 buf.append(line) 

72 if len(buf) >= _NDJSON_BATCH_SIZE: 

73 yield "".join(buf) 

74 buf.clear() 

75 if buf: 

76 yield "".join(buf) 

77 

78 

79@task_instances_log_router.get( 

80 "/{task_id}/logs/{try_number}", 

81 responses={ 

82 **create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

83 status.HTTP_200_OK: { 

84 "description": "Successful Response", 

85 "content": ndjson_example_response_for_get_log, 

86 }, 

87 }, 

88 dependencies=[Depends(requires_access_dag("GET", DagAccessEntity.TASK_LOGS))], 

89 response_model=TaskInstancesLogResponse, 

90 response_model_exclude_unset=True, 

91) 

92def get_log( 

93 dag_id: str, 

94 dag_run_id: str, 

95 task_id: str, 

96 try_number: NonNegativeInt, 

97 accept: HeaderAcceptJsonOrNdjson, 

98 request: Request, 

99 dag_bag: DagBagDep, 

100 session: SessionDep, 

101 full_content: bool = False, 

102 map_index: int = -1, 

103 token: str | None = None, 

104): 

105 """Get logs for a specific task instance.""" 

106 if not token: 

107 metadata = {} 

108 else: 

109 try: 

110 metadata = URLSafeSerializer(request.app.state.secret_key).loads(token) 

111 except BadSignature: 

112 raise HTTPException( 

113 status.HTTP_400_BAD_REQUEST, "Bad Signature. Please use only the tokens provided by the API." 

114 ) 

115 

116 if metadata.get("download_logs") and metadata["download_logs"]: 116 ↛ 117line 116 didn't jump to line 117 because the condition on line 116 was never true

117 full_content = True 

118 

119 metadata["download_logs"] = full_content 

120 

121 task_log_reader = TaskLogReader() 

122 

123 if not task_log_reader.supports_read: 123 ↛ 124line 123 didn't jump to line 124 because the condition on line 123 was never true

124 raise HTTPException(status.HTTP_400_BAD_REQUEST, "Task log handler does not support read logs.") 

125 

126 query = ( 

127 select(TaskInstance) 

128 .where( 

129 TaskInstance.task_id == task_id, 

130 TaskInstance.dag_id == dag_id, 

131 TaskInstance.run_id == dag_run_id, 

132 TaskInstance.map_index == map_index, 

133 TaskInstance.try_number == try_number, 

134 ) 

135 .join(TaskInstance.dag_run) 

136 .options(joinedload(TaskInstance.trigger).joinedload(Trigger.triggerer_job)) 

137 .options(joinedload(TaskInstance.dag_model)) 

138 ) 

139 ti = session.scalar(query) 

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

141 query = ( 

142 select(TaskInstanceHistory) 

143 .where( 

144 TaskInstanceHistory.task_id == task_id, 

145 TaskInstanceHistory.dag_id == dag_id, 

146 TaskInstanceHistory.run_id == dag_run_id, 

147 TaskInstanceHistory.map_index == map_index, 

148 TaskInstanceHistory.try_number == try_number, 

149 ) 

150 .options(joinedload(TaskInstanceHistory.dag_run)) 

151 # we need to joinedload the dag_run, since FileTaskHandler._render_filename needs ti.dag_run 

152 ) 

153 ti = session.scalar(query) 

154 

155 if ti is None: 

156 metadata["end_of_log"] = True 

157 raise HTTPException(status.HTTP_404_NOT_FOUND, "TaskInstance not found") 

158 

159 dag = dag_bag.get_dag_for_run(ti.dag_run, session=session) 

160 if dag: 160 ↛ 164line 160 didn't jump to line 164 because the condition on line 160 was always true

161 with contextlib.suppress(TaskNotFound): 

162 ti.task = dag.get_task(ti.task_id) 

163 

164 if accept == Mimetype.NDJSON: # only specified application/x-ndjson will return streaming response 164 ↛ 166line 164 didn't jump to line 166 because the condition on line 164 was never true

165 # LogMetadata(TypedDict) is used as type annotation for log_reader; added ignore to suppress mypy error 

166 raw_stream = task_log_reader.read_log_stream(ti, try_number, metadata) # type: ignore[arg-type] 

167 log_stream = _buffered_ndjson_stream(raw_stream) 

168 headers = None 

169 if not metadata.get("end_of_log", False): 

170 headers = { 

171 "Airflow-Continuation-Token": URLSafeSerializer(request.app.state.secret_key).dumps(metadata) 

172 } 

173 return StreamingResponse(media_type="application/x-ndjson", content=log_stream, headers=headers) 

174 

175 # application/json, or something else we don't understand. 

176 # Return JSON format, which will be more easily for users to debug. 

177 

178 # LogMetadata(TypedDict) is used as type annotation for log_reader; added ignore to suppress mypy error 

179 structured_log_stream, out_metadata = task_log_reader.read_log_chunks(ti, try_number, metadata) # type: ignore[arg-type] 

180 encoded_token = None 

181 if not out_metadata.get("end_of_log", False): 181 ↛ 182line 181 didn't jump to line 182 because the condition on line 181 was never true

182 encoded_token = URLSafeSerializer(request.app.state.secret_key).dumps(out_metadata) 

183 return TaskInstancesLogResponse.model_construct( 

184 continuation_token=encoded_token, content=list(structured_log_stream) 

185 ) 

186 

187 

188@task_instances_log_router.get( 

189 "/{task_id}/externalLogUrl/{try_number}", 

190 responses=create_openapi_http_exception_doc([status.HTTP_400_BAD_REQUEST, status.HTTP_404_NOT_FOUND]), 

191 dependencies=[Depends(requires_access_dag("GET", DagAccessEntity.TASK_LOGS))], 

192) 

193def get_external_log_url( 

194 dag_id: str, 

195 dag_run_id: str, 

196 task_id: str, 

197 try_number: PositiveInt, 

198 session: SessionDep, 

199 map_index: int = -1, 

200) -> ExternalLogUrlResponse: 

201 """Get external log URL for a specific task instance.""" 

202 task_log_reader = TaskLogReader() 

203 

204 if not task_log_reader.supports_external_link: 204 ↛ 208line 204 didn't jump to line 208 because the condition on line 204 was always true

205 raise HTTPException(status.HTTP_400_BAD_REQUEST, "Task log handler does not support external logs.") 

206 

207 # Fetch the task instance 

208 query = ( 

209 select(TaskInstance) 

210 .where( 

211 TaskInstance.task_id == task_id, 

212 TaskInstance.dag_id == dag_id, 

213 TaskInstance.run_id == dag_run_id, 

214 TaskInstance.map_index == map_index, 

215 ) 

216 .options(joinedload(TaskInstance.dag_model)) 

217 ) 

218 ti = session.scalar(query) 

219 

220 if ti is None: 

221 raise HTTPException(status.HTTP_404_NOT_FOUND, "TaskInstance not found") 

222 

223 url = task_log_reader.log_handler.get_external_log_url(ti, try_number) 

224 return ExternalLogUrlResponse(url=url)