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

101 statements  

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

1""" 

2Routes for interacting with work queue objects. 

3""" 

4 

5from typing import List, Optional 

6from uuid import UUID 

7 

8from fastapi import ( 

9 Body, 

10 Depends, 

11 Header, 

12 HTTPException, 

13 Path, 

14 status, 

15) 

16from sqlalchemy.exc import IntegrityError 

17 

18import prefect.server.api.dependencies as dependencies 

19import prefect.server.models as models 

20import prefect.server.schemas as schemas 

21from prefect.server.database import ( 

22 PrefectDBInterface, 

23 provide_database_interface, 

24) 

25from prefect.server.models.deployments import mark_deployments_ready 

26from prefect.server.models.work_queues import ( 

27 emit_work_queue_status_event, 

28 mark_work_queues_ready, 

29) 

30from prefect.server.schemas.statuses import WorkQueueStatus 

31from prefect.server.utilities.server import PrefectRouter 

32from prefect.types import DateTime 

33 

34router: PrefectRouter = PrefectRouter(prefix="/work_queues", tags=["Work Queues"]) 

35 

36 

37@router.post("/", status_code=status.HTTP_201_CREATED) 

38async def create_work_queue( 

39 work_queue: schemas.actions.WorkQueueCreate, 

40 db: PrefectDBInterface = Depends(provide_database_interface), 

41) -> schemas.responses.WorkQueueResponse: 

42 """ 

43 Creates a new work queue. 

44 

45 If a work queue with the same name already exists, an error 

46 will be raised. 

47 

48 For more information, see https://docs.prefect.io/v3/concepts/work-pools#work-queues. 

49 """ 

50 

51 try: 

52 async with db.session_context(begin_transaction=True) as session: 

53 model = await models.work_queues.create_work_queue( 

54 session=session, work_queue=work_queue 

55 ) 

56 

57 response = schemas.responses.WorkQueueResponse.model_validate( 

58 model, from_attributes=True 

59 ) 

60 if response.concurrency_limit is not None: 

61 response.active_slots = 0 

62 except IntegrityError: 

63 raise HTTPException( 

64 status_code=status.HTTP_409_CONFLICT, 

65 detail="A work queue with this name already exists.", 

66 ) 

67 

68 return response 

69 

70 

71@router.patch("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT) 

72async def update_work_queue( 

73 work_queue: schemas.actions.WorkQueueUpdate, 

74 work_queue_id: UUID = Path(..., description="The work queue id", alias="id"), 

75 db: PrefectDBInterface = Depends(provide_database_interface), 

76) -> None: 

77 """ 

78 Updates an existing work queue. 

79 """ 

80 async with db.session_context(begin_transaction=True) as session: 

81 existing_work_queue = await models.work_queues.read_work_queue( 

82 session=session, work_queue_id=work_queue_id 

83 ) 

84 if existing_work_queue is None: 

85 raise HTTPException( 

86 status_code=status.HTTP_404_NOT_FOUND, 

87 detail=f"Work Queue {work_queue_id} not found", 

88 ) 

89 

90 result = await models.work_queues.update_work_queue( 

91 session=session, 

92 work_queue_id=work_queue_id, 

93 work_queue=work_queue, 

94 emit_status_change=emit_work_queue_status_event, 

95 ) 

96 if not result: 96 ↛ 97line 96 didn't jump to line 97 because the condition on line 96 was never true

97 raise HTTPException( 

98 status_code=status.HTTP_404_NOT_FOUND, 

99 detail=f"Work Queue {work_queue_id} not found", 

100 ) 

101 

102 

103@router.get("/name/{name}") 

104async def read_work_queue_by_name( 

105 name: str = Path(..., description="The work queue name"), 

106 db: PrefectDBInterface = Depends(provide_database_interface), 

107) -> schemas.responses.WorkQueueResponse: 

108 """ 

109 Get a work queue by id. 

110 """ 

111 async with db.session_context() as session: 

112 work_queue = await models.work_queues.read_work_queue_by_name( 

113 session=session, name=name 

114 ) 

115 if not work_queue: 

116 raise HTTPException( 

117 status_code=status.HTTP_404_NOT_FOUND, detail="work queue not found" 

118 ) 

119 response = schemas.responses.WorkQueueResponse.model_validate( 

120 work_queue, from_attributes=True 

121 ) 

122 if response.concurrency_limit is not None: 

123 response.active_slots = ( 

124 await models.work_queues.count_work_queue_active_slots( 

125 session=session, work_queue_id=work_queue.id 

126 ) 

127 ) 

128 return response 

129 

130 

131@router.get("/{id:uuid}") 

132async def read_work_queue( 

133 work_queue_id: UUID = Path(..., description="The work queue id", alias="id"), 

134 db: PrefectDBInterface = Depends(provide_database_interface), 

135) -> schemas.responses.WorkQueueResponse: 

136 """ 

137 Get a work queue by id. 

138 """ 

139 async with db.session_context() as session: 

140 work_queue = await models.work_queues.read_work_queue( 

141 session=session, work_queue_id=work_queue_id 

142 ) 

143 if not work_queue: 

144 raise HTTPException( 

145 status_code=status.HTTP_404_NOT_FOUND, detail="work queue not found" 

146 ) 

147 response = schemas.responses.WorkQueueResponse.model_validate( 

148 work_queue, from_attributes=True 

149 ) 

150 if response.concurrency_limit is not None: 150 ↛ anywhereline 150 didn't jump anywhere: it always raised an exception.

151 response.active_slots = ( 

152 await models.work_queues.count_work_queue_active_slots( 

153 session=session, work_queue_id=work_queue_id 

154 ) 

155 ) 

156 return response 

157 

158 

159@router.post("/{id:uuid}/get_runs") 

160async def read_work_queue_runs( 

161 docket: dependencies.Docket, 

162 work_queue_id: UUID = Path(..., description="The work queue id", alias="id"), 

163 limit: int = dependencies.LimitBody(), 

164 scheduled_before: DateTime = Body( 

165 None, 

166 description=( 

167 "Only flow runs scheduled to start before this time will be returned." 

168 ), 

169 ), 

170 x_prefect_ui: Optional[bool] = Header( 

171 default=False, 

172 description="A header to indicate this request came from the Prefect UI.", 

173 ), 

174 db: PrefectDBInterface = Depends(provide_database_interface), 

175) -> List[schemas.responses.FlowRunResponse]: 

176 """ 

177 Get flow runs from the work queue. 

178 """ 

179 async with db.session_context(begin_transaction=True) as session: 

180 work_queue, flow_runs = await models.work_queues.get_runs_in_work_queue( 

181 session=session, 

182 work_queue_id=work_queue_id, 

183 scheduled_before=scheduled_before, 

184 limit=limit, 

185 ) 

186 

187 # The Prefect UI often calls this route to see which runs are enqueued. 

188 # We do not want to record this as an actual poll event. 

189 if x_prefect_ui: 189 ↛ 190line 189 didn't jump to line 190 because the condition on line 189 was never true

190 return flow_runs 

191 

192 await docket.add( 

193 mark_work_queues_ready, 

194 key=f"mark_work_queues_ready:{work_queue_id}", 

195 )( 

196 polled_work_queue_ids=[work_queue_id], 

197 ready_work_queue_ids=( 

198 [work_queue_id] if work_queue.status == WorkQueueStatus.NOT_READY else [] 

199 ), 

200 ) 

201 

202 await docket.add( 

203 mark_deployments_ready, 

204 key=f"mark_deployments_ready:work_queue:{work_queue_id}", 

205 )( 

206 work_queue_ids=[work_queue_id], 

207 ) 

208 

209 return flow_runs 

210 

211 

212@router.post("/filter") 

213async def read_work_queues( 

214 limit: int = dependencies.LimitBody(), 

215 offset: int = Body(0, ge=0), 

216 work_queues: Optional[schemas.filters.WorkQueueFilter] = None, 

217 db: PrefectDBInterface = Depends(provide_database_interface), 

218) -> List[schemas.responses.WorkQueueResponse]: 

219 """ 

220 Query for work queues. 

221 """ 

222 async with db.session_context() as session: 

223 wqs = await models.work_queues.read_work_queues( 

224 session=session, offset=offset, limit=limit, work_queue_filter=work_queues 

225 ) 

226 

227 ret = [ 

228 schemas.responses.WorkQueueResponse.model_validate(wq, from_attributes=True) 

229 for wq in wqs 

230 ] 

231 queues_with_limit = [wq for wq in ret if wq.concurrency_limit is not None] 

232 if queues_with_limit: 

233 slot_counts = await models.work_queues.count_work_queue_active_slots_bulk( 

234 session=session, 

235 work_queue_ids=[wq.id for wq in queues_with_limit], 

236 ) 

237 for wq_response in queues_with_limit: 

238 wq_response.active_slots = slot_counts.get(wq_response.id, 0) 

239 

240 return ret 

241 

242 

243@router.delete("/{id:uuid}", status_code=status.HTTP_204_NO_CONTENT) 

244async def delete_work_queue( 

245 work_queue_id: UUID = Path(..., description="The work queue id", alias="id"), 

246 db: PrefectDBInterface = Depends(provide_database_interface), 

247) -> None: 

248 """ 

249 Delete a work queue by id. 

250 """ 

251 async with db.session_context(begin_transaction=True) as session: 

252 existing_work_queue = await models.work_queues.read_work_queue( 

253 session=session, work_queue_id=work_queue_id 

254 ) 

255 if existing_work_queue is None: 

256 raise HTTPException( 

257 status_code=status.HTTP_404_NOT_FOUND, detail="work queue not found" 

258 ) 

259 result = await models.work_queues.delete_work_queue( 

260 session=session, work_queue_id=work_queue_id 

261 ) 

262 if not result: 

263 raise HTTPException( 

264 status_code=status.HTTP_404_NOT_FOUND, detail="work queue not found" 

265 ) 

266 

267 

268@router.post("/{id:uuid}/concurrency_status") 

269async def read_work_queue_concurrency_status( 

270 work_queue_id: UUID = Path(..., description="The work queue id", alias="id"), 

271 page: int = Body(1, ge=1), 

272 limit: int = dependencies.LimitBody(), 

273 db: PrefectDBInterface = Depends(provide_database_interface), 

274) -> schemas.responses.WorkQueueConcurrencyStatus: 

275 """ 

276 Read concurrency status for a work queue, including paginated flow run 

277 summaries. active_slots always reflects the total count. 

278 """ 

279 import asyncio 

280 

281 from prefect.types._datetime import now as prefect_now 

282 

283 run_offset = (page - 1) * limit 

284 

285 async with db.session_context() as session: 

286 work_queue = await models.work_queues.read_work_queue( 

287 session=session, work_queue_id=work_queue_id 

288 ) 

289 if not work_queue: 

290 raise HTTPException( 

291 status_code=status.HTTP_404_NOT_FOUND, 

292 detail="Work queue not found.", 

293 ) 

294 

295 # Count and paginated fetch in parallel 

296 total_count, slot_holders_page = await asyncio.gather( 

297 models.workers.count_work_queue_slot_holders( 

298 session=session, work_queue_id=work_queue_id 

299 ), 

300 models.workers.get_work_queue_slot_holders( 

301 session=session, 

302 work_queue_id=work_queue_id, 

303 offset=run_offset, 

304 limit=limit, 

305 ), 

306 ) 

307 

308 current_time = prefect_now("UTC") 

309 

310 flow_runs = [ 

311 schemas.responses.FlowRunSlotSummary( 

312 id=run.id, 

313 name=run.name, 

314 state_type=run.state_type if run.state_type else None, 

315 state_name=run.state_name if run.state_name else None, 

316 start_time=run.start_time, 

317 state_timestamp=run.state_timestamp, 

318 time_in_current_state=( 

319 (current_time - run.state_timestamp) if run.state_timestamp else None 

320 ), 

321 ) 

322 for run, slot_acquired_at in slot_holders_page 

323 ] 

324 

325 return schemas.responses.WorkQueueConcurrencyStatus( 

326 active_slots=total_count, 

327 concurrency_limit=work_queue.concurrency_limit, 

328 flow_runs=flow_runs, 

329 count=total_count, 

330 limit=limit, 

331 pages=(total_count + limit - 1) // limit if limit > 0 else 0, 

332 page=page, 

333 ) 

334 

335 

336@router.get("/{id:uuid}/status") 

337async def read_work_queue_status( 

338 work_queue_id: UUID = Path(..., description="The work queue id", alias="id"), 

339 db: PrefectDBInterface = Depends(provide_database_interface), 

340) -> schemas.core.WorkQueueStatusDetail: 

341 """ 

342 Get the status of a work queue. 

343 """ 

344 async with db.session_context() as session: 

345 work_queue_status = await models.work_queues.read_work_queue_status( 

346 session=session, work_queue_id=work_queue_id 

347 ) 

348 return work_queue_status