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
« 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.
18from __future__ import annotations
20import contextlib
21import textwrap
22from collections.abc import Generator, Iterable
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
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
45_NDJSON_BATCH_SIZE = conf.getint("api", "log_stream_buffer_size")
47task_instances_log_router = AirflowRouter(
48 tags=["Task Instance"], prefix="/dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances"
49)
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}
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)
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 )
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
119 metadata["download_logs"] = full_content
121 task_log_reader = TaskLogReader()
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.")
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)
155 if ti is None:
156 metadata["end_of_log"] = True
157 raise HTTPException(status.HTTP_404_NOT_FOUND, "TaskInstance not found")
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)
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)
175 # application/json, or something else we don't understand.
176 # Return JSON format, which will be more easily for users to debug.
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 )
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()
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.")
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)
220 if ti is None:
221 raise HTTPException(status.HTTP_404_NOT_FOUND, "TaskInstance not found")
223 url = task_log_reader.log_handler.get_external_log_url(ti, try_number)
224 return ExternalLogUrlResponse(url=url)