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
« 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"""
5from typing import List, Optional
6from uuid import UUID
8from fastapi import (
9 Body,
10 Depends,
11 Header,
12 HTTPException,
13 Path,
14 status,
15)
16from sqlalchemy.exc import IntegrityError
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
34router: PrefectRouter = PrefectRouter(prefix="/work_queues", tags=["Work Queues"])
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.
45 If a work queue with the same name already exists, an error
46 will be raised.
48 For more information, see https://docs.prefect.io/v3/concepts/work-pools#work-queues.
49 """
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 )
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 )
68 return response
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 )
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 )
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
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
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 )
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
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 )
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 )
209 return flow_runs
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 )
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)
240 return ret
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 )
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
281 from prefect.types._datetime import now as prefect_now
283 run_offset = (page - 1) * limit
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 )
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 )
308 current_time = prefect_now("UTC")
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 ]
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 )
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