Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/api/events.py: 44%

122 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1import base64 

2from typing import TYPE_CHECKING, List, Optional 

3 

4from fastapi import Response, WebSocket 

5from fastapi.exceptions import HTTPException 

6from fastapi.param_functions import Depends, Path 

7from fastapi.params import Body, Query 

8from sqlalchemy.ext.asyncio import AsyncSession 

9from starlette.requests import Request 

10from starlette.status import WS_1002_PROTOCOL_ERROR 

11 

12from prefect._internal.compatibility.starlette import status 

13from prefect.logging import get_logger 

14from prefect.server.api.dependencies import is_ephemeral_request 

15from prefect.server.database import PrefectDBInterface, provide_database_interface 

16from prefect.server.events import messaging, stream 

17from prefect.server.events.counting import ( 

18 Countable, 

19 InvalidEventCountParameters, 

20 TimeUnit, 

21) 

22from prefect.server.events.filters import EventFilter, EventOrder 

23from prefect.server.events.models.automations import automations_session 

24from prefect.server.events.pipeline import EventsPipeline 

25from prefect.server.events.schemas.events import Event, EventCount, EventPage 

26from prefect.server.events.storage import ( 

27 INTERACTIVE_PAGE_SIZE, 

28 InvalidTokenError, 

29 database, 

30) 

31from prefect.server.utilities import subscriptions 

32from prefect.server.utilities.server import PrefectRouter 

33from prefect.settings import ( 

34 PREFECT_EVENTS_MAXIMUM_WEBSOCKET_BACKFILL, 

35 PREFECT_EVENTS_WEBSOCKET_BACKFILL_PAGE_SIZE, 

36) 

37 

38if TYPE_CHECKING: 38 ↛ 39line 38 didn't jump to line 39 because the condition on line 38 was never true

39 import logging 

40 

41logger: "logging.Logger" = get_logger(__name__) 

42 

43 

44router: PrefectRouter = PrefectRouter(prefix="/events", tags=["Events"]) 

45 

46 

47@router.post("", status_code=status.HTTP_204_NO_CONTENT, response_class=Response) 

48async def create_events( 

49 events: List[Event], 

50 ephemeral_request: bool = Depends(is_ephemeral_request), 

51) -> None: 

52 """ 

53 Record a batch of Events. 

54 

55 For more information, see https://docs.prefect.io/v3/concepts/events. 

56 """ 

57 if ephemeral_request: 57 ↛ 58line 57 didn't jump to line 58 because the condition on line 57 was never true

58 await EventsPipeline().process_events(events) 

59 else: 

60 received_events = [event.receive() for event in events] 

61 await messaging.publish(received_events) 

62 

63 

64@router.websocket("/in") 

65async def stream_events_in(websocket: WebSocket) -> None: 

66 """Open a WebSocket to stream incoming Events""" 

67 websocket = await subscriptions.accept_prefect_socket(websocket) 

68 if not websocket: 

69 return 

70 

71 try: 

72 async with messaging.create_event_publisher() as publisher: 

73 async for event_json in websocket.iter_text(): 

74 event = Event.model_validate_json(event_json) 

75 await publisher.publish_event(event.receive()) 

76 except subscriptions.NORMAL_DISCONNECT_EXCEPTIONS: # pragma: no cover 

77 pass # it's fine if a client disconnects either normally or abnormally 

78 

79 return None 

80 

81 

82@router.websocket("/out") 

83async def stream_workspace_events_out( 

84 websocket: WebSocket, 

85) -> None: 

86 """Open a WebSocket to stream Events""" 

87 websocket = await subscriptions.accept_prefect_socket( 

88 websocket, 

89 ) 

90 if not websocket: 

91 return 

92 

93 try: 

94 # After authentication, the next message is expected to be a filter message, any 

95 # other type of message will close the connection. 

96 message = await websocket.receive_json() 

97 

98 if message["type"] != "filter": 

99 return await websocket.close( 

100 WS_1002_PROTOCOL_ERROR, reason="Expected 'filter' message" 

101 ) 

102 

103 wants_backfill = message.get("backfill", True) 

104 

105 try: 

106 filter = EventFilter.model_validate(message["filter"]) 

107 except Exception as e: 

108 return await websocket.close( 

109 WS_1002_PROTOCOL_ERROR, reason=f"Invalid filter: {e}" 

110 ) 

111 

112 filter.occurred.clamp(PREFECT_EVENTS_MAXIMUM_WEBSOCKET_BACKFILL.value()) 

113 filter.order = EventOrder.ASC 

114 

115 # subscribe to the ongoing event stream first so we don't miss events... 

116 async with stream.events(filter) as event_stream: 

117 # ...then if the user wants, backfill up to the last 1k events... 

118 if wants_backfill: 

119 backfilled_ids = set() 

120 

121 async with automations_session() as session: 

122 backfill, _, next_page = await database.query_events( 

123 session=session, 

124 filter=filter, 

125 page_size=PREFECT_EVENTS_WEBSOCKET_BACKFILL_PAGE_SIZE.value(), 

126 ) 

127 

128 while backfill: 

129 for event in backfill: 

130 backfilled_ids.add(event.id) 

131 await websocket.send_json( 

132 { 

133 "type": "event", 

134 "event": event.model_dump(mode="json"), 

135 } 

136 ) 

137 

138 if not next_page: 

139 break 

140 

141 backfill, _, next_page = await database.query_next_page( 

142 session=session, 

143 page_token=next_page, 

144 ) 

145 

146 # ...before resuming the ongoing stream of events 

147 async for event in event_stream: 

148 if not event: 

149 if await subscriptions.still_connected(websocket): 

150 continue 

151 break 

152 

153 if wants_backfill and event.id in backfilled_ids: 

154 backfilled_ids.remove(event.id) 

155 continue 

156 

157 await websocket.send_json( 

158 {"type": "event", "event": event.model_dump(mode="json")} 

159 ) 

160 

161 except subscriptions.NORMAL_DISCONNECT_EXCEPTIONS: # pragma: no cover 

162 pass # it's fine if a client disconnects either normally or abnormally 

163 

164 return None 

165 

166 

167def verified_page_token( 

168 page_token: str = Query(..., alias="page-token"), 

169) -> str: 

170 try: 

171 page_token = base64.b64decode(page_token.encode()).decode() 

172 except Exception: 

173 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN) 

174 

175 if not page_token: 175 ↛ 178line 175 didn't jump to line 178 because the condition on line 175 was always true

176 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN) 

177 

178 return page_token 

179 

180 

181@router.post( 

182 "/filter", 

183) 

184async def read_events( 

185 request: Request, 

186 filter: Optional[EventFilter] = Body( 

187 None, 

188 description=( 

189 "Additional optional filter criteria to narrow down the set of Events" 

190 ), 

191 ), 

192 limit: int = Body( 

193 INTERACTIVE_PAGE_SIZE, 

194 ge=0, 

195 le=INTERACTIVE_PAGE_SIZE, 

196 embed=True, 

197 description="The number of events to return with each page", 

198 ), 

199 db: PrefectDBInterface = Depends(provide_database_interface), 

200) -> EventPage: 

201 """ 

202 Queries for Events matching the given filter criteria in the given Account. Returns 

203 the first page of results, and the URL to request the next page (if there are more 

204 results). 

205 """ 

206 filter = filter or EventFilter() 

207 async with db.session_context() as session: 

208 events, total, next_token = await database.query_events( 

209 session=session, 

210 filter=filter, 

211 page_size=limit, 

212 ) 

213 

214 return EventPage( 

215 events=events, 

216 total=total, 

217 next_page=generate_next_page_link(request, next_token), 

218 ) 

219 

220 

221@router.get( 

222 "/filter/next", 

223) 

224async def read_account_events_page( 

225 request: Request, 

226 page_token: str = Depends(verified_page_token), 

227 db: PrefectDBInterface = Depends(provide_database_interface), 

228) -> EventPage: 

229 """ 

230 Returns the next page of Events for a previous query against the given Account, and 

231 the URL to request the next page (if there are more results). 

232 """ 

233 async with db.session_context() as session: 

234 try: 

235 events, total, next_token = await database.query_next_page( 

236 session=session, page_token=page_token 

237 ) 

238 except InvalidTokenError: 

239 raise HTTPException(status_code=status.HTTP_403_FORBIDDEN) 

240 

241 return EventPage( 

242 events=events, 

243 total=total, 

244 next_page=generate_next_page_link(request, next_token), 

245 ) 

246 

247 

248def generate_next_page_link( 

249 request: Request, 

250 page_token: Optional[str], 

251) -> Optional[str]: 

252 if not page_token: 

253 return None 

254 

255 next_page = ( 

256 f"{request.base_url}api/events/filter/next" 

257 f"?page-token={base64.b64encode(page_token.encode()).decode()}" 

258 ) 

259 return next_page 

260 

261 

262@router.post( 

263 "/count-by/{countable}", 

264) 

265async def count_account_events( 

266 filter: EventFilter, 

267 countable: Countable = Path(...), 

268 time_unit: TimeUnit = Body(default=TimeUnit.day), 

269 time_interval: float = Body(default=1.0, ge=0.01), 

270 db: PrefectDBInterface = Depends(provide_database_interface), 

271) -> List[EventCount]: 

272 """ 

273 Returns distinct objects and the count of events associated with them. Objects 

274 that can be counted include the day the event occurred, the type of event, or 

275 the IDs of the resources associated with the event. 

276 """ 

277 async with db.session_context() as session: 

278 return await handle_event_count_request( 

279 session=session, 

280 filter=filter, 

281 countable=countable, 

282 time_unit=time_unit, 

283 time_interval=time_interval, 

284 ) 

285 

286 

287async def handle_event_count_request( 

288 session: AsyncSession, 

289 filter: EventFilter, 

290 countable: Countable, 

291 time_unit: TimeUnit, 

292 time_interval: float, 

293) -> List[EventCount]: 

294 try: 

295 return await database.count_events( 

296 session=session, 

297 filter=filter, 

298 countable=countable, 

299 time_unit=time_unit, 

300 time_interval=time_interval, 

301 ) 

302 except InvalidEventCountParameters as exc: 

303 raise HTTPException( 

304 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, 

305 detail=exc.message, 

306 )