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

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 datetime import datetime 

20from typing import Annotated 

21 

22from fastapi import Depends, HTTPException, status 

23from sqlalchemy import select 

24from sqlalchemy.orm import joinedload 

25 

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 

53 

54event_logs_router = AirflowRouter(tags=["Event Log"], prefix="/eventLogs") 

55 

56 

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 

80 

81 

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) 

212 

213 return EventLogCollectionResponse( 

214 event_logs=event_logs, 

215 total_entries=total_entries, 

216 )