Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/event_logs.py: 100%
26 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
19from datetime import datetime
20from typing import Annotated
22from fastapi import Depends, HTTPException, status
23from sqlalchemy import select
24from sqlalchemy.orm import joinedload
26from airflow.api_fastapi.common.db.common import (
27 SessionDep,
28 paginated_select,
29)
30from airflow.api_fastapi.common.parameters import (
31 FilterOptionEnum,
32 FilterParam,
33 QueryLimit,
34 QueryOffset,
35 SortParam,
36 _PrefixSearchParam,
37 _SearchParam,
38 filter_param_factory,
39 prefix_search_param_factory,
40 search_param_factory,
41)
42from airflow.api_fastapi.common.router import AirflowRouter
43from airflow.api_fastapi.core_api.datamodels.event_logs import (
44 EventLogCollectionResponse,
45 EventLogResponse,
46)
47from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
48from airflow.api_fastapi.core_api.security import (
49 ReadableEventLogsFilterDep,
50 requires_access_event_log,
51)
52from airflow.models import Log
54event_logs_router = AirflowRouter(tags=["Event Log"], prefix="/eventLogs")
57@event_logs_router.get(
58 "/{event_log_id}",
59 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
60 dependencies=[Depends(requires_access_event_log("GET"))],
61)
62def get_event_log(
63 event_log_id: int,
64 session: SessionDep,
65) -> EventLogResponse:
66 event_log = session.scalar(
67 # Log.dttm is nullable at the DB level, but EventLogResponse.when is a non-optional
68 # datetime. Rows with dttm=NULL would cause a Pydantic validation error (500), so
69 # exclude them here. Such rows can exist in legacy installs or via direct DB inserts
70 # that bypass Log.__init__ (which always sets dttm = timezone.utcnow()).
71 # Making EventLogResponse.when nullable would be a breaking API contract change for
72 # clients that currently rely on `when` always being present.
73 select(Log)
74 .where(Log.id == event_log_id, Log.dttm.is_not(None))
75 .options(joinedload(Log.task_instance))
76 )
77 if event_log is None:
78 raise HTTPException(status.HTTP_404_NOT_FOUND, f"The Event Log with id: `{event_log_id}` not found")
79 return event_log
82@event_logs_router.get(
83 "",
84 dependencies=[Depends(requires_access_event_log("GET"))],
85)
86def get_event_logs(
87 limit: QueryLimit,
88 offset: QueryOffset,
89 session: SessionDep,
90 order_by: Annotated[
91 SortParam,
92 Depends(
93 SortParam(
94 [
95 "id", # event_log_id
96 "dttm", # when
97 "dag_id",
98 "task_id",
99 "run_id",
100 "event",
101 "logical_date",
102 "owner",
103 "extra",
104 ],
105 Log,
106 to_replace={"when": "dttm", "event_log_id": "id"},
107 ).dynamic_depends()
108 ),
109 ],
110 # Exact match filters (for backward compatibility)
111 dag_id: Annotated[FilterParam[str | None], Depends(filter_param_factory(Log.dag_id, str | None))],
112 task_id: Annotated[FilterParam[str | None], Depends(filter_param_factory(Log.task_id, str | None))],
113 run_id: Annotated[FilterParam[str | None], Depends(filter_param_factory(Log.run_id, str | None))],
114 map_index: Annotated[FilterParam[int | None], Depends(filter_param_factory(Log.map_index, int | None))],
115 try_number: Annotated[FilterParam[int | None], Depends(filter_param_factory(Log.try_number, int | None))],
116 owner: Annotated[FilterParam[str | None], Depends(filter_param_factory(Log.owner, str | None))],
117 event: Annotated[FilterParam[str | None], Depends(filter_param_factory(Log.event, str | None))],
118 excluded_events: Annotated[
119 FilterParam[list[str] | None],
120 Depends(
121 filter_param_factory(Log.event, list[str] | None, FilterOptionEnum.NOT_IN, "excluded_events")
122 ),
123 ],
124 included_events: Annotated[
125 FilterParam[list[str] | None],
126 Depends(filter_param_factory(Log.event, list[str] | None, FilterOptionEnum.IN, "included_events")),
127 ],
128 before: Annotated[
129 FilterParam[datetime | None],
130 Depends(filter_param_factory(Log.dttm, datetime | None, FilterOptionEnum.LESS_THAN, "before")),
131 ],
132 after: Annotated[
133 FilterParam[datetime | None],
134 Depends(filter_param_factory(Log.dttm, datetime | None, FilterOptionEnum.GREATER_THAN, "after")),
135 ],
136 # Pattern search filters (substring match, ILIKE)
137 dag_id_pattern: Annotated[_SearchParam, Depends(search_param_factory(Log.dag_id, "dag_id_pattern"))],
138 task_id_pattern: Annotated[_SearchParam, Depends(search_param_factory(Log.task_id, "task_id_pattern"))],
139 run_id_pattern: Annotated[_SearchParam, Depends(search_param_factory(Log.run_id, "run_id_pattern"))],
140 owner_pattern: Annotated[_SearchParam, Depends(search_param_factory(Log.owner, "owner_pattern"))],
141 event_pattern: Annotated[_SearchParam, Depends(search_param_factory(Log.event, "event_pattern"))],
142 # Prefix pattern search filters (index-friendly, case-sensitive)
143 dag_id_prefix_pattern: Annotated[
144 _PrefixSearchParam,
145 Depends(prefix_search_param_factory(Log.dag_id, "dag_id_prefix_pattern")),
146 ],
147 task_id_prefix_pattern: Annotated[
148 _PrefixSearchParam,
149 Depends(prefix_search_param_factory(Log.task_id, "task_id_prefix_pattern")),
150 ],
151 run_id_prefix_pattern: Annotated[
152 _PrefixSearchParam,
153 Depends(prefix_search_param_factory(Log.run_id, "run_id_prefix_pattern")),
154 ],
155 owner_prefix_pattern: Annotated[
156 _PrefixSearchParam,
157 Depends(prefix_search_param_factory(Log.owner, "owner_prefix_pattern")),
158 ],
159 event_prefix_pattern: Annotated[
160 _PrefixSearchParam,
161 Depends(prefix_search_param_factory(Log.event, "event_prefix_pattern")),
162 ],
163 readable_event_logs_filter: ReadableEventLogsFilterDep,
164) -> EventLogCollectionResponse:
165 """Get all Event Logs."""
166 query = (
167 # Log.dttm is nullable at the DB level, but EventLogResponse.when is a non-optional
168 # datetime. Rows with dttm=NULL would cause a Pydantic validation error (500), so
169 # exclude them here. Such rows can exist in legacy installs or via direct DB inserts
170 # that bypass Log.__init__ (which always sets dttm = timezone.utcnow()).
171 # Making EventLogResponse.when nullable would be a breaking API contract change for
172 # clients that currently rely on `when` always being present.
173 select(Log)
174 .where(Log.dttm.is_not(None))
175 .options(joinedload(Log.task_instance), joinedload(Log.dag_model))
176 )
177 event_logs_select, total_entries = paginated_select(
178 statement=query,
179 order_by=order_by,
180 filters=[
181 # Exact match filters
182 dag_id,
183 task_id,
184 run_id,
185 map_index,
186 try_number,
187 owner,
188 event,
189 excluded_events,
190 included_events,
191 before,
192 after,
193 # Pattern search filters
194 dag_id_pattern,
195 dag_id_prefix_pattern,
196 task_id_pattern,
197 task_id_prefix_pattern,
198 run_id_pattern,
199 run_id_prefix_pattern,
200 owner_pattern,
201 owner_prefix_pattern,
202 event_pattern,
203 event_prefix_pattern,
204 # Permission
205 readable_event_logs_filter,
206 ],
207 offset=offset,
208 limit=limit,
209 session=session,
210 )
211 event_logs = session.scalars(event_logs_select)
213 return EventLogCollectionResponse(
214 event_logs=event_logs,
215 total_entries=total_entries,
216 )